每次迭代返回一行(不含换行符) * * @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 { // 步骤 1:clearstatcache + 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; } } // 步骤 4:fseek 到 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); // 步骤 7:maxLines 达到即保存并退出 if ($yielded >= $maxLines) { $newOffset = $bufferOffset + $pos; $this->savePosition($filePath, $inode, $newOffset, $lastLineHash); return; } } // 步骤 6:buffer 中已消费 $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; } } } // 步骤 8:EOF 处理 if ($buffer === '') { // 文件恰好在最后一个 \n 处结束:offset = 文件末尾(ftell) $finalOffset = ftell($fp); $this->savePosition($filePath, $inode, (int) $finalOffset, $lastLineHash); } else { // 尾部半行(无 \n):可能是日志正在写入,不 yield,offset 回退到半行起点,下次重读 $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, ]); } }