大文件处理管道:日志分析与背压

本节目标

真实运维里你常要分析几十 GB 的日志:统计错误数、按分钟聚合、过滤 WARN。把整文件读进来 split('\n') 根本不现实。正确姿势是管道(pipeline):源头可读流 → 若干转换流 → 终点可写流,pipeline 负责把错误一路传播、统一收尾。

// 运行环境:Node.js 14+(自动生成约 20MB 日志演示,无需联网)
// 保存为 buf-l5.js,执行:node buf-l5.js
const fs = require('fs')
const path = require('path')
const { Transform } = require('stream')
const { pipeline } = require('stream/promises')

const logFile = path.join(__dirname, 'app.log')
// 造一份 20MB 混合日志,含 INFO / WARN / ERROR
const lines = []
for (let i = 0; i < 200000; i++) {
  const lvl = i % 7 === 0 ? 'ERROR' : (i % 3 === 0 ? 'WARN' : 'INFO')
  lines.push(`2026-08-12 10:${String(i % 60).padStart(2, '0')}:00 [${lvl}] request #${i}`)
}
fs.writeFileSync(logFile, lines.join('\n'))

let errorCount = 0
// 转换流:每收到一块,就地统计其中的 ERROR 数量(不缓存整文件)
const counter = new Transform({
  transform(chunk, enc, cb) {
    const errs = (chunk.toString().match(/ERROR/g) || []).length
    errorCount += errs
    cb(null, chunk) // 原样转发出去(这里只是顺便计数)
  }
})

async function main() {
  // pipeline 会自动串联并传播错误;任意一环出错都会正确销毁所有流
  await pipeline(
    fs.createReadStream(logFile),
    counter,
    fs.createWriteStream(path.join(__dirname, 'app-filtered.log'))
  )
  console.log('ERROR 行数       :', errorCount)
  console.log('全程堆内存峰值约 :', (process.memoryUsage().heapUsed / 1024 / 1024).toFixed(1), 'MB')
  console.log('(文件 20MB,内存却没有随文件大小暴涨)')
}
main().catch((e) => { console.error(e); process.exit(1) })
     .finally(() => {
       fs.unlinkSync(logFile)
       fs.unlinkSync(path.join(__dirname, 'app-filtered.log'))
     })

背压在这一节如何体现? pipeline 底层用 pipe 连接各环节。当 counter 下游的写文件较慢时,write 会返回 false,pipeline 自动让上游暂停产数据;文件写完腾出空间,'drain' 触发后继续。于是内存恒定在"缓冲区 + 当前块",不会因文件 20MB 就吃 20MB 内存。

忽略背压的代价:如果你不用 pipe/pipeline,而是 src.on('data', c => ws.write(c)),且不处理 ws.write 返回 false 时该暂停——数据会无限堆积,文件越大内存爆得越狠。这就是"背压失控"。

名词解释

课后练习

  1. 把 counter 里的 cb(null, chunk) 改成 cb()(不转发数据),会发生什么?
    • 答案:下游写文件那环将收不到任何数据,输出文件为空。cb 的第二个参数是"交给下游的块",省略就等于产出空块。统计仍会执行(errorCount 照常)。
  2. 为什么统计 ERROR 要在"块"上用正则,而不是 chunk.toString().split('\n') 逐行?
    • 答案:一块的边界可能刚好切断某行(一行被拆在两个 chunk),split('\n') 会漏判半行;用 match(/ERROR/g) 数出现次数不受行边界影响。若要精确按行,需自己维护"跨块的行缓冲"。
  3. stream.pipeline 相比手动 a.pipe(b).pipe(c) 强在哪?
    • 答案:pipeline 会监听所有流的 error 并统一销毁整条链、在 await 时抛错、处理完成后回调;手动 pipe 链若中间流报错,容易出现资源未释放或错误被吞。生产代码一律用 pipeline。

总结

pipeline 是流的"工程化形态"——它把你之前手动拼的 pipe 升级成了一条带错误处理、资源自动回收的流水线。我在这节刻意用 20MB 文件、却看到内存几乎不动,就是要让你直观感受流的威力:处理的成本只和"当前这块"有关,与"文件总大小"解耦。但也请记住背压的另一面教训——pipeline 之所以稳,是因为它内部替你管好了背压;一旦你退回到手动 on('data') 且忽略 write 的返回值,那个"内存恒定"的魔法就消失了,取而代之的是悄无声息的内存膨胀。所以结论很明确:生产环境里,凡是流的串联,优先 pipeline,不要手写裸 pipe。把你今天写的 counter 换成"按分钟聚合""过滤 WARN 写另一份文件",就是真实日志平台的雏形——而它和动辄几个 G 的日志打交道时,内存依旧纹丝不动。