每次迭代返回一行(键值与 DebugLogToolkit::FIELDS 一致) * * @throws \RuntimeException 文件无法打开或 stat/fseek/fread 失败时抛出(交由上层 try-catch) */ public function read(string $filePath, int $maxLines = 100000): \Generator { // per-run 状态复位(同一个 reader 实例可串行处理多个文件) $this->reachedEof = false; $this->droppedLines = 0; $this->currentLineHash = null; $this->parseFails = 0; $this->failSamples = []; $fp = @fopen($filePath, 'rb'); if ($fp === false) { throw new \RuntimeException("DebugLogReaderService: 无法打开文件 {$filePath}"); } try { clearstatcache(true, $filePath); $stat = fstat($fp); if ($stat === false) { throw new \RuntimeException("DebugLogReaderService: fstat 失败 {$filePath}"); } $currentInode = (int) $stat['ino']; $currentSize = (int) $stat['size']; $this->currentInode = $currentInode; // 读 position(无则默认 offset=0) $offset = 0; $position = $this->getLastPosition($filePath); if ($position !== null) { $offset = (int) $position['offset']; // 轮转检测:inode 变化(同日文件被重建)OR size < offset(防御性归零) if ((int) $position['inode'] !== $currentInode || $currentSize < $offset) { $offset = 0; } } $this->currentOffset = $offset; if ($offset > 0 && fseek($fp, $offset) !== 0) { throw new \RuntimeException("DebugLogReaderService: fseek 失败 offset={$offset} file={$filePath}"); } $yielded = 0; $buffer = ''; // $bufferOffset 始终等于 buffer[0] 对应的文件字节偏移 $bufferOffset = $offset; while (true) { $chunk = fread($fp, self::BUFFER_SIZE); if ($chunk === false) { throw new \RuntimeException("DebugLogReaderService: fread 失败 file={$filePath}"); } if ($chunk === '') { // EOF:buffer 为空说明干净收尾;有半行则 currentOffset 已停在最后完整行末尾 $this->reachedEof = ($buffer === ''); return; } $buffer .= $chunk; // 内层循环:处理 buffer 中所有完整行(含末尾 \n 的) $pos = 0; while (($nlPos = strpos($buffer, "\n", $pos)) !== false) { $line = substr($buffer, $pos, $nlPos - $pos); $pos = $nlPos + 1; // 先推进安全偏移再产出/判失败:坏行也算已消费(毒行不重读), // 消费方在事务内 savePosition 拿到的 offset 与已处理行严格对齐 $this->currentOffset = $bufferOffset + $pos; $row = DebugLogToolkit::decodeLine($line); if ($row === null) { $this->parseFails++; $this->recordFailSample($line); continue; } $this->currentLineHash = md5($line); yield $row; $yielded++; if ($yielded >= $maxLines) { return; } } // buffer 中已消费 $pos 字节,半行(buffer[$pos..])保留 $bufferOffset += $pos; $buffer = substr($buffer, $pos); // 超长行防御:整行丢弃(禁止强制切行——会把一条 JSON 切成两条坏行) if (strlen($buffer) >= self::MAX_LINE_BYTES) { $this->droppedLines++; $buffer = ''; // 扫描至下一个换行符,将其后内容作为新 buffer 继续 while (true) { $chunk = fread($fp, self::BUFFER_SIZE); if ($chunk === false || $chunk === '') { // 至 EOF 仍未等到换行符:整个尾部按已丢弃处理,偏移推到文件末尾 $bufferOffset = (int) ftell($fp); $this->currentOffset = $bufferOffset; $this->reachedEof = true; return; } $nl = strpos($chunk, "\n"); if ($nl !== false) { $bufferOffset = (int) ftell($fp) - (strlen($chunk) - $nl - 1); $this->currentOffset = $bufferOffset; $buffer = substr($chunk, $nl + 1); break; } } } } } finally { if (is_resource($fp)) { fclose($fp); } } } /** * 增量读取遗留 CSV 文件,yield 关联数组化的数据行. * * 读取协议(与 read() 的差异全因 CSV 的多物理行特性): * 1. fopen 'rb' + fstat 取 inode/size + position 轮转检测归零——与 read() 一致 * 2. fseek 到 offset 后循环 fgetcsv():fgetcsv 是有状态解析,原生处理引号内 * 嵌入换行(一条逻辑记录跨多物理行),返回值即一条完整逻辑记录 * 【禁止】按 \n 切行后 str_getcsv——引号内换行会被撕碎 * 3. 每读出一条逻辑记录(无论产出/表头跳过/失败)都把 currentOffset 推进到 * ftell($fp)(记录末尾),保证事务内 savePosition 与已消费记录对齐 * 4. position==0(从文件头开始)时首行经 DebugLogToolkit::isCsvHeader() 判定, * 是表头则跳过不产出(offset 仍推进);续读 offset 处不会是表头,不判定 * 5. 行字段数 !== 8 计 parseFails 不产出(失败样本进 failSamples) * 6. 行字段按 DebugLogToolkit::FIELDS 顺序 map 成关联数组;create_time 为 * 数字字符串时转 int,create_time_title 保留字符串(CSV 一切皆字符串) * 7. 产出 maxRows 条后退出(不自动 savePosition,由上层事务内落盘) * * 注:CSV 路径无"尾部半行"概念——遗留 CSV 是已封存的静态文件(非追加中), * fgetcsv 读到 EOF 返 false 即干净收尾(reachedEof=true)。 * * @param string $filePath csv 文件绝对路径 * @param int $maxRows 单次产出的最大数据行数(表头/失败行不计) * * @return \Generator 每次迭代返回一行(键值与 DebugLogToolkit::FIELDS 一致) * * @throws \RuntimeException 文件无法打开或 stat/fseek 失败时抛出(交由上层 try-catch) */ public function readCsv(string $filePath, int $maxRows = 100000): \Generator { // per-run 状态复位(与 read() 同一套状态位,CSV 路径 droppedLines 恒 0) $this->reachedEof = false; $this->droppedLines = 0; $this->currentLineHash = null; $this->parseFails = 0; $this->failSamples = []; $fp = @fopen($filePath, 'rb'); if ($fp === false) { throw new \RuntimeException("DebugLogReaderService: 无法打开文件 {$filePath}"); } try { clearstatcache(true, $filePath); $stat = fstat($fp); if ($stat === false) { throw new \RuntimeException("DebugLogReaderService: fstat 失败 {$filePath}"); } $currentInode = (int) $stat['ino']; $currentSize = (int) $stat['size']; $this->currentInode = $currentInode; // 读 position(无则默认 offset=0) $offset = 0; $position = $this->getLastPosition($filePath); if ($position !== null) { $offset = (int) $position['offset']; // 轮转检测:inode 变化 OR size < offset(防御性归零) if ((int) $position['inode'] !== $currentInode || $currentSize < $offset) { $offset = 0; } } $this->currentOffset = $offset; if ($offset > 0 && fseek($fp, $offset) !== 0) { throw new \RuntimeException("DebugLogReaderService: fseek 失败 offset={$offset} file={$filePath}"); } // 仅从文件头(offset==0)开始读时才做表头判定:续读起点不会是表头 $checkHeader = ($offset === 0); $yielded = 0; while (true) { // 显式传全部参数:separator/enclosure/escape 与遗留 fputcsv 写出参数一致 // (PHP 8.4+ 省略 escape 会触发 deprecation,显式传可前向兼容) $fields = fgetcsv($fp, null, ',', '"', '\\'); if ($fields === false) { // fgetcsv 返 false 即 EOF(干净收尾,无半行概念) $this->reachedEof = true; return; } // 每读出一条逻辑记录就把安全偏移推进到记录末尾(含引号内换行占的 // 多物理行)——表头/失败行同样推进(毒行不重读),保证事务内 // savePosition 与已消费记录严格对齐 $this->currentOffset = (int) ftell($fp); if ($checkHeader) { $checkHeader = false; if (DebugLogToolkit::isCsvHeader($this->normalizeCsvFields($fields))) { // 表头行:跳过不产出(offset 已推进) continue; } } if (count($fields) !== 8) { $this->parseFails++; $this->recordFailSample('csv-fields=' . count($fields) . ': ' . implode(',', $this->stringifyCsvFields($fields))); continue; } $row = []; foreach (DebugLogToolkit::FIELDS as $i => $field) { $value = $fields[$i] ?? ''; // fgetcsv 对空字段可能返回 null,统一落为 ''(与 JSONL 行的空串语义对齐) $row[$field] = $value === null ? '' : $value; } // create_time 数字字符串转 int(时间戳恒为非负整数);title 保留字符串 if (ctype_digit((string) $row['create_time'])) { $row['create_time'] = (int) $row['create_time']; } $this->currentLineHash = md5(serialize($fields)); yield $row; $yielded++; if ($yielded >= $maxRows) { return; } } } finally { if (is_resource($fp)) { fclose($fp); } } } /** * 当前安全读取位置(read/readCsv 期间/结束后可调). * * @return array{inode: int, offset: int, last_line_hash: string|null} */ public function getCurrentPosition(): array { return [ 'inode' => $this->currentInode, 'offset' => $this->currentOffset, 'last_line_hash' => $this->currentLineHash, ]; } /** * 本轮读取是否干净到达 EOF(JSONL:尾部无半行;CSV:fgetcsv 返 false). * * 导入侧据此判断"文件已导完":reachedEof 且 offset >= filesize 才允许 unlink。 */ public function reachedEof(): bool { return $this->reachedEof; } /** * 本轮 read 丢弃的超长行数(由导入侧并入 fails 统计;CSV 路径恒 0). */ public function getDroppedLines(): int { return $this->droppedLines; } /** * 本轮读取解析失败计数(JSONL 坏行 + CSV 字段数不符,由导入侧传给 savePosition). */ public function getParseFails(): int { return $this->parseFails; } /** * 解析失败样本(json 数组字符串或 null,由导入侧传给 savePosition). * * 每条样本截断至 FAIL_SAMPLE_LENGTH 字符,最多 MAX_FAIL_SAMPLES 条。 */ public function getFailSamples(): ?string { if ($this->failSamples === []) { return null; } $json = json_encode($this->failSamples, JSON_UNESCAPED_UNICODE); return $json === false ? null : $json; } /** * 查询文件的上次读取位置(按 node_id 隔离). * * @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 = DebugLogPosition::where('file_path', $filePath) ->where('node_id', $nodeId) ->find(); 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'), ]; } /** * 保存读取位置(upsert,按 node_id 隔离). * * 由导入任务在"insertBatch 同一事务"内显式调用(本类不自动保存), * 保证崩溃时入库与进度要么同时生效、要么同时回滚,消重复导入窗口。 * * @param string $filePath 日志文件绝对路径 * @param int $inode 当前文件 inode * @param int $offset 新的字节偏移(最后一条已消费记录末尾) * @param string|null $lastLineHash 最后一条已产出记录 md5 * @param int $failCount 本批解析失败累计次数 * @param string|null $failSamples 解析失败样本(json 或 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 = DebugLogPosition::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; } DebugLogPosition::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, ]); } /** * 记录一条解析失败样本(截断至 FAIL_SAMPLE_LENGTH,超量丢弃). */ protected function recordFailSample(string $raw): void { if (count($this->failSamples) >= self::MAX_FAIL_SAMPLES) { return; } $this->failSamples[] = mb_substr($raw, 0, self::FAIL_SAMPLE_LENGTH); } /** * CSV 字段数组归一化后交 isCsvHeader 判定(null 字段落为 ''). * * @param array $fields fgetcsv 原始返回值 */ protected function normalizeCsvFields(array $fields): array { return array_map(static function ($value) { return $value === null ? '' : $value; }, $fields); } /** * CSV 字段数组字符串化(失败样本拼接用,null 落为空串占位). * * @param array $fields fgetcsv 原始返回值 * * @return string[] */ protected function stringifyCsvFields(array $fields): array { return array_map(static function ($value) { return $value === null ? '' : (string) $value; }, $fields); } }