Files
ulthon_admin/extend/base/common/service/NginxLogReaderServiceBase.php
augushong 8faa52ce90 feat(nginx-log): 多节点适配(node_id 全链路 + 全局/按节点聚合 + 节点筛选 + v2.4.0 升级脚本)
- 6 张表加 node_id 字段 + position/stat 改复合唯一键(migration)
- 6 个 Scheme 同步 node_id 注解 + 唯一键
- Reader 带 node_id 参数 + 懒加载接管(getLastPosition 回退查 node_id='')
- Aggregator 全局行(node_id='')+ 按节点循环 insertAggregatesForScope
- Stat 所有查询方法加 where node_id 条件
- Dashboard 控制器/视图加节点筛选下拉框 + AJAX 带 node_id
- AccessLog 列表加节点列 + 采集进度卡片显示节点信息
- import 定时任务 run_type 改 all(每节点采集自己的日志)
- v2.4.0 升级脚本追加 ALTER TABLE(幂等 check-then-alter)
- 菜单调整:去掉 Nginx 顶级菜单 + 读取位置管理,改为系统管理下挂两个子菜单
- ulthon-timer 文档补充 run_type=all 多节点部署注意事项
2026-08-08 10:33:20 +08:00

282 lines
12 KiB
PHP
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

<?php
namespace base\common\service;
use app\admin\model\NginxLogPosition;
use app\common\service\HostService;
/**
* Nginx 日志增量读取 serviceBase 层).
*
* 职责:
* - 按文件 offset 增量读取日志行Generator不读整个文件到内存
* - 检测日志轮转inode 变化 OR filesize < offset覆盖 rename+rebuild 与 copytruncate 两种模式)
* - 将读取进度持久化到 NginxLogPosition 表upsert按 node_id 隔离)
*
* 多节点适配position 表按 (node_id, file_path) 唯一定位读取进度。
* 旧数据node_id 为空)通过懒加载接管(见 getLastPosition首个查询到它的本节点原子 UPDATE node_id
* 避免 migration 全表回填(多节点共享 DB 时 migration 只跑一次,只有懒加载能确保每节点各自接管)。
*
* 不负责解析T5 Parser、聚合T7 Aggregator
*
* 依赖倒置:内部 model 调用走 `app\admin\model\NginxLogPosition`App 入口类),
* 使用者可在 app/ 重写该 model 拦截行为;本类不直接 `new NginxLogPosition`(静态调用满足多态)。
* 节点身份走 `app\common\service\HostService`App 入口类)。
*/
class NginxLogReaderServiceBase
{
/** 单次 fread 的 buffer 大小64KB */
protected const BUFFER_SIZE = 65536;
/** 单行无换行符时的硬上限8MB超过则强制按一行 yield防止 OOM */
protected const MAX_LINE_BYTES = 8388608;
/**
* 增量读取日志文件yield 完整行.
*
* 读取协议:
* 1. fopen + fstat 获取当前 inode 和 filesize
* 2. 读 position 行无则视为新建inode=当前、offset=0
* 3. 轮转检测双条件inode 不一致 OR filesize < position.offset → 重置 offset=0
* 4. fseek 到 offset
* 5. 每读 64KB buffer定位已读段内最后 `\n`,仅 yield 该位置之前的完整行
* 6. 新 offset = 最后 `\n` + 1buffer 中剩余半行保留到下次循环)
* 7. yield maxLines 后提前 return 并 savePosition
* 8. EOF 时 savePositionoffset = 文件末尾(若尾部有半行未 yieldoffset 回退到半行起点)
*
* @param string $filePath 日志文件绝对路径
* @param int $maxLines 单次 yield 的最大行数(达到即提前退出并保存 offset
*
* @return \Generator<string> 每次迭代返回一行(不含换行符)
*
* @throws \RuntimeException 文件无法打开或 stat/fseek 失败时抛出(交由上层 try-catch
*/
public function read(string $filePath, int $maxLines = 100000): \Generator
{
$fp = @fopen($filePath, 'rb');
if ($fp === false) {
throw new \RuntimeException("NginxLogReaderService: 无法打开文件 {$filePath}");
}
try {
// 步骤 1clearstatcache + fstat 取当前 inode/size
clearstatcache(true, $filePath);
$stat = fstat($fp);
if ($stat === false) {
throw new \RuntimeException("NginxLogReaderService: fstat 失败 {$filePath}");
}
$currentInode = (int) $stat['ino'];
$currentSize = (int) $stat['size'];
// 步骤 2读 position无则默认 offset=0
$position = $this->getLastPosition($filePath);
$inode = $currentInode;
$offset = 0;
if ($position !== null) {
$inode = (int) $position['inode'];
$offset = (int) $position['offset'];
// 步骤 3轮转检测双条件
// inode 不一致 → rename+rebuild 模式Linux inode 有效时主判定)
// filesize < offset → copytruncate 模式 / Windows 上 inode 恒为 0 时的主判定
if ($inode !== $currentInode || $currentSize < $offset) {
$offset = 0;
$inode = $currentInode;
}
}
// 步骤 4fseek 到 offset
if ($offset > 0 && fseek($fp, $offset) !== 0) {
throw new \RuntimeException("NginxLogReaderService: fseek 失败 offset={$offset} file={$filePath}");
}
$yielded = 0;
$lastLineHash = null;
$buffer = '';
// $bufferOffset 始终等于 buffer[0] 对应的文件字节偏移
$bufferOffset = $offset;
// 步骤 5-7循环读取 + yield
while (true) {
$chunk = fread($fp, self::BUFFER_SIZE);
if ($chunk === false) {
// 读错误:交上层处理,不保存进度(下次重读)
throw new \RuntimeException("NginxLogReaderService: fread 失败 file={$filePath}");
}
if ($chunk === '') {
// EOF
break;
}
$buffer .= $chunk;
// 内层循环:处理 buffer 中所有完整行(含末尾 \n 的)
$pos = 0;
while (($nlPos = strpos($buffer, "\n", $pos)) !== false) {
$line = substr($buffer, $pos, $nlPos - $pos);
$pos = $nlPos + 1;
yield $line;
$yielded++;
$lastLineHash = md5($line);
// 步骤 7maxLines 达到即保存并退出
if ($yielded >= $maxLines) {
$newOffset = $bufferOffset + $pos;
$this->savePosition($filePath, $inode, $newOffset, $lastLineHash);
return;
}
}
// 步骤 6buffer 中已消费 $pos 字节半行buffer[$pos..])保留
$bufferOffset += $pos;
$buffer = substr($buffer, $pos);
// 防御buffer 无 \n 累积过长(异常畸形行),强制按一行处理避免 OOM
if (strlen($buffer) >= self::MAX_LINE_BYTES) {
yield $buffer;
$yielded++;
$lastLineHash = md5($buffer);
$bufferOffset += strlen($buffer);
$buffer = '';
if ($yielded >= $maxLines) {
$this->savePosition($filePath, $inode, $bufferOffset, $lastLineHash);
return;
}
}
}
// 步骤 8EOF 处理
if ($buffer === '') {
// 文件恰好在最后一个 \n 处结束offset = 文件末尾ftell
$finalOffset = ftell($fp);
$this->savePosition($filePath, $inode, (int) $finalOffset, $lastLineHash);
} else {
// 尾部半行(无 \n可能是日志正在写入不 yieldoffset 回退到半行起点,下次重读
$this->savePosition($filePath, $inode, $bufferOffset, $lastLineHash);
}
} finally {
if (is_resource($fp)) {
fclose($fp);
}
}
}
/**
* 查询文件的上次读取位置.
*
* 多节点适配where 条件带 node_id 隔离各节点进度。
* 懒加载接管:多节点共享 DB 时 migration 只跑一次,旧记录 node_id 为空,
* 首个查询到它的本节点原子 UPDATE node_id 接管,保证每节点各自有独立 position 记录。
*
* @param string $filePath 日志文件绝对路径
*
* @return array|null 命中时返回 [inode, offset, last_line_hash, parse_fail_count, parse_fail_samples],无记录返回 null
*/
public function getLastPosition(string $filePath): ?array
{
$nodeId = HostService::getNodeId();
$row = NginxLogPosition::where('file_path', $filePath)
->where('node_id', $nodeId)
->find();
if ($row === null) {
// 懒加载接管:旧记录 node_id 为空,原子抢占接管到当前节点
$row = $this->lazyAdoptPosition($filePath, $nodeId);
if ($row === null) {
return null;
}
}
return [
'inode' => (int) $row->getData('inode'),
'offset' => (int) $row->getData('offset'),
'last_line_hash' => $row->getData('last_line_hash'),
'parse_fail_count' => (int) $row->getData('parse_fail_count'),
'parse_fail_samples' => $row->getData('parse_fail_samples'),
];
}
/**
* 懒加载接管:把旧的无 node_id 记录原子接管到当前节点.
*
* 多节点共享 DB 时 migration 只跑一次,无法为每节点预置 position 记录。
* 旧记录node_id='')由首个查询到它的本节点通过原子条件 UPDATE 接管:
* UPDATE ... SET node_id=current WHERE node_id='' AND file_path=X
* 并发安全WHERE node_id='' 保证只有一个节点能 affected=1其他节点查询时已被接管。
*
* @return \think\Model|null 接管成功返回接管后的记录,无旧记录或被其他节点抢占返回 null
*/
protected function lazyAdoptPosition(string $filePath, string $nodeId)
{
$legacy = NginxLogPosition::where('file_path', $filePath)
->where('node_id', '')
->find();
if ($legacy === null) {
return null;
}
// 原子条件 UPDATE只有 node_id='' 时才能被接管(防并发多节点同时接管同一记录)
$affected = NginxLogPosition::where('id', $legacy->getData('id'))
->where('node_id', '')
->update(['node_id' => $nodeId]);
if ($affected === 0) {
// 被其他节点抢先接管,本节点放弃(下次 savePosition 会创建新记录)
return null;
}
return NginxLogPosition::where('file_path', $filePath)
->where('node_id', $nodeId)
->find();
}
/**
* 保存读取位置upsert.
*
* Reader 内部在 maxLines 达到或 EOF 时调用,传 lastLineHash最后一行 md5
* failCount 默认 0每次新读取批次开始时由 Reader 重置;上层若需记录解析失败,
* 在 parse 后再次调用本方法覆盖 failCount/failSamples 即可offset 保持一致)。
*
* @param string $filePath 日志文件绝对路径(主键)
* @param int $inode 当前文件 inode
* @param int $offset 新的字节偏移
* @param string|null $lastLineHash 最后一行 md5用于去重/检测)
* @param int $failCount 解析失败累计次数(默认 0
* @param string|null $failSamples 解析失败样本(默认 null
*/
public function savePosition(string $filePath, int $inode, int $offset, ?string $lastLineHash, int $failCount = 0, ?string $failSamples = null): void
{
$now = time();
$nodeId = HostService::getNodeId();
$existing = NginxLogPosition::where('file_path', $filePath)
->where('node_id', $nodeId)
->find();
if ($existing !== null) {
$existing->save([
'inode' => $inode,
'offset' => $offset,
'last_read_time' => $now,
'last_line_hash' => $lastLineHash,
'parse_fail_count' => $failCount,
'parse_fail_samples' => $failSamples,
]);
return;
}
NginxLogPosition::create([
'node_id' => $nodeId,
'file_path' => $filePath,
'inode' => $inode,
'offset' => $offset,
'last_read_time' => $now,
'last_line_hash' => $lastLineHash,
'parse_fail_count' => $failCount,
'parse_fail_samples' => $failSamples,
'create_time' => $now,
'update_time' => $now,
]);
}
}