人人都会AI编程

25.2 流的高级应用:自定义流、转换流、压缩流

更新时间:2026-07-10

在第 7 章中我们已经系统学习过 Node.js 四种流类型的基本概念和背压机制。但在实际项目中,仅仅使用内置的可读/可写流往往不够——我们经常需要处理自定义的数据转换、实现自定义协议的解析,或者与压缩、加密算法无缝集成。本节将重点介绍如何编写生产级的自定义流,以及如何利用压缩流为系统提升效率。

25.2.1 自定义可读流:从一个范例看 Readable 的实现

Node.js 的 stream.Readable 是创建可读流的基类。我们只需要继承它并实现 _read(size) 方法,内部通过 this.push(chunk) 向消费者提供数据。当数据推送完毕时,推入 null 表示流结束。

下面是一个生成斐波那契数列的可读流示例,它不会一次性把所有数字算完,而是按需生成,天然支持背压:

const { Readable } = require('stream');

class FibonacciStream extends Readable {
  constructor(max) {
    super({ objectMode: true }); // 允许推送非 Buffer 数据
    this.max = max;
    this.n = 0;
    this.a = 0;
    this.b = 1;
  }

  _read(size) {
    if (this.n >= this.max) {
      this.push(null); // 数据已全部推送
      return;
    }

    const value = this.a;
    [this.a, this.b] = [this.b, this.a + this.b];
    this.n++;

    // push 返回 false 表示消费者读取较慢,内部缓冲已满,应当暂停生产
    const shouldContinue = this.push(value);
    if (!shouldContinue) {
      // 暂停生成,等待消费者调用 _read 再次触发
      return;
    }
    // 如果 push 返回 true,可以继续同步推送,但一般不会循环调用 _read
    // 实际中通常配合异步源使用,这里不加额外逻辑
  }
}

// 使用
const fib = new FibonacciStream(10);
fib.on('data', (num) => console.log(num));
fib.on('end', () => console.log('结束'));

关键点解析:

  • objectMode: true 让流处理 JavaScript 对象而非 Buffer,适用于结构化数据。
  • _read(size) 被消费者触发,size 是一个建议值,流可以忽略。
  • push() 返回 false 是背压信号:内部缓冲区已超过 highWaterMark,表明消费者处理速度跟不上,我们应暂停数据生成。下次消费者读取数据时,_read 会被再次调用。
  • 在实际异步数据源(如文件读取、数据库游标)中,我们通常在 _read 中拉取一批数据,并在异步回调中 push,利用 push 的返回值决定是否继续拉取。

25.2.2 自定义可写流:高效处理批量写入

stream.Writable 是自定义可写流的基类,必须实现 _write(chunk, encoding, callback) 方法。每次写入时,该方法被调用,处理完数据后调用 callback 通知流已完成此块处理,可以继续接收下一个数据块。还可以可选实现 _writev(chunks, callback) 来批量处理多个缓冲块,减少系统调用。

下面是一个简单的“行收集器”,将流式数据按行缓存,最后一次性输出所有行的例子:

const { Writable } = require('stream');

class LineCollector extends Writable {
  constructor(options) {
    super({ ...options, objectMode: true }); // 兼容字符串模式
    this.lines = [];
  }

  _write(chunk, encoding, callback) {
    // chunk 可能是 Buffer 或字符串
    const str = chunk.toString();
    this.lines.push(...str.split('\n').filter(line => line !== ''));
    callback(); // 无需异步操作,直接回调
  }

  _final(callback) {
    // 所有数据写入完毕时触发
    console.log('收集到的行:', this.lines);
    callback();
  }
}

// 使用
const collector = new LineCollector();
collector.write('hello\nworld\n');
collector.write('foo\nbar');
collector.end();

背压与 _writev 优化:
在高吞吐量场景下,如果每次 _write 都触发一次回调,可能产生较大的函数调用开销。实现 _writev 可以将多个连续写入的 chunk 累积成数组一次性处理,类似于批量写入数据库:

class BatchWriter extends Writable {
  constructor(options) {
    super({ ...options, highWaterMark: 16 });
    // 注意:_writev 仅在缓冲中存在多个待处理块时被调用
  }

  _writev(chunks, callback) {
    const data = chunks.map(c => c.chunk.toString()).join('');
    console.log('批量写入:', data);
    // 模拟异步
    setTimeout(callback, 10);
  }

  _write(chunk, encoding, callback) {
    // 可选降级,当只有一个块时走 _write
    console.log('单个写入:', chunk.toString());
    callback();
  }
}

_writev 对于需要顺序写入的数据库或网络套接字尤为有用,可以减少 I/O 次数。

25.2.3 自定义转换流:Transform 的巧妙运用

stream.Transform 继承自 Duplex,它的精髓在于在读写之间插入一个数据转换逻辑。需要实现 _transform(chunk, encoding, callback) 方法,每接收一个数据块,可以多次调用 this.push(transformedChunk) 将转换后的数据输出到可读端。当所有输入结束时,_flush(callback) 可以用来输出剩余数据(例如编码器剩余的缓冲区)。

示例 1:JSON 行解析器
将接收到的字节流按行分割,并解析每一行为 JSON 对象:

const { Transform } = require('stream');

class JSONLineParser extends Transform {
  constructor() {
    super({ objectMode: true }); // 输入为 Buffer,输出为对象
    this.buffer = '';
  }

  _transform(chunk, encoding, callback) {
    this.buffer += chunk.toString();
    const lines = this.buffer.split('\n');
    // 最后一行可能不完整,保留到下次处理
    this.buffer = lines.pop();
    for (const line of lines) {
      if (line.trim() === '') continue;
      try {
        this.push(JSON.parse(line));
      } catch (err) {
        // 解析失败触发错误,流会销毁
        callback(err);
        return;
      }
    }
    callback();
  }

  _flush(callback) {
    if (this.buffer.trim()) {
      try {
        this.push(JSON.parse(this.buffer));
      } catch (err) {
        callback(err);
        return;
      }
    }
    callback();
  }
}

// 使用
const parser = new JSONLineParser();
parser.on('data', obj => console.log('解析对象:', obj));
parser.write('{"name":"Alice"}\n{"name":"Bob"}\n');
parser.end('{"name":"Charlie"}');

示例 2:流式加密/解密(Cipher 转换流)
Node.js 的 crypto 模块提供了针对流的 Cipher/Decipher 类,它们本身也是 Transform 流。我们也可以基于 Transform 自定义简单的 XOR 加密演示:

class XORCipher extends Transform {
  constructor(key) {
    super();
    this.key = Buffer.from(key);
    this.keyIndex = 0;
  }

  _transform(chunk, encoding, callback) {
    const result = Buffer.alloc(chunk.length);
    for (let i = 0; i < chunk.length; i++) {
      result[i] = chunk[i] ^ this.key[this.keyIndex++ % this.key.length];
    }
    this.push(result);
    callback();
  }
}

25.2.4 压缩流:无缝集成 zlib 的数据管道

Node.js 内置的 zlib 模块提供了一系列压缩/解压缩 API,并且每一组算法都有对应的流式接口:createGzipcreateGunzipcreateDeflatecreateBrotliCompress 等。这些方法返回的都是 Transform 流,可以直接通过 pipe 串联到管道中。

场景一:HTTP 响应压缩
在 Web 服务器中按需压缩响应体可以大幅减少带宽占用。以 Koa 或 Express 为例,我们可以手动创建一个压缩中间件:

const zlib = require('zlib');
const { Transform } = require('stream');

function gzipMiddleware(req, res, next) {
  const acceptEncoding = req.headers['accept-encoding'] || '';
  if (!acceptEncoding.includes('gzip')) {
    return next();
  }

  // 替换 res 的写方法,让响应经过 gzip 压缩后发送
  const gzip = zlib.createGzip();
  res.setHeader('Content-Encoding', 'gzip');
  
  // 原始 res 的写入方法需要被替换,但直接操作原生 http 模块较复杂
  // 实际可使用框架的压缩中间件,此处仅为演示原理
  // 省略完整实现...
}

更常用的方式是借助成熟的中间件(如 compression),但其背后就是 zlib 的 Transform 流。

场景二:文件压缩备份
需求:将一个大日志文件压缩后写入备份文件,同时计算压缩后文件的 MD5 值。这可以利用管道同时进行压缩和哈希计算:

const fs = require('fs');
const zlib = require('zlib');
const crypto = require('crypto');

const readStream = fs.createReadStream('access.log');
const writeStream = fs.createWriteStream('access.log.gz');
const gzip = zlib.createGzip();
const md5 = crypto.createHash('md5', { encoding: 'hex' });

// 使用 pipeline 自动处理背压和错误
const { pipeline } = require('stream');
pipeline(
  readStream,
  gzip,
  md5, // 哈希流也是 Transform
  writeStream,
  (err) => {
    if (err) {
      console.error('压缩失败', err);
    } else {
      console.log('压缩完成,MD5:', md5.digest('hex'));
    }
  }
);

这里注意:crypto.createHash 返回的也是一个 Transform 流,它接收数据并在最后输出摘要,但在管道中需要调用 digest() 来获取结果。由于 pipeline 结束后流已关闭,所以可以在回调中安全使用。

场景三:压缩/解压转换流的封装
有时我们需要一个既能压缩又能解压的通用流,可以根据参数动态切换:

function createZipStream(format = 'gzip') {
  switch (format) {
    case 'gzip': return zlib.createGzip();
    case 'deflate': return zlib.createDeflate();
    case 'brotli': return zlib.createBrotliCompress();
    default: throw new Error('Unsupported format');
  }
}

25.2.5 组合多个转换流:管道的力量

Node.js 流的一大优势在于可组合性。通过 pipe 或者 stream.pipeline,我们可以将自定义转换流与内置流(压缩、加密、文件读写)串联起来,形成清晰的数据处理管道。例如:

读取文件 → 解密 → 解压 → 自定义解析 → 写入数据库

这样的管道不仅逻辑清晰,而且自动处理背压和错误传播,是处理大规模数据的经典模式。

注意事项:

  • 错误处理:为管道中每个流添加 error 监听,或者使用 pipeline 统一捕获错误。一旦某个流出错,管道会自动销毁未结束的流。
  • 内存控制:使用 objectMode 时,highWaterMark 表示允许缓冲的对象个数,而非字节数。处理大量小对象可能导致内存膨胀,需要合理设置阈值。
  • 结束与清理:自定义流应在 _final_destroy 中释放资源(如关闭文件描述符、清除定时器)。

25.2.6 综合实例:实时日志收集与压缩传输

假设我们要开发一个客户端,监控某个日志文件的新增内容,实时压缩后通过网络发送到服务端。利用自定义可读流监听文件变化,再通过压缩流和 socket 管道可以轻松实现:

const { createReadStream } = require('fs');
const { createGzip } = require('zlib');
const net = require('net');
const { pipeline } = require('stream');

// 以 tail 方式监听文件,略去具体实现
const logTail = new LogTailStream('/var/log/app.log');

const client = net.createConnection({ port: 5000 }, () => {
  pipeline(
    logTail,           // 自定义可读流,推送新增行
    createGzip(),      // 压缩
    client,            // 网络套接字
    (err) => {
      if (err) console.error('管道失败', err);
    }
  );
});

这一套架构充分利用了自定义流和压缩流的组合,实现了高效、低延迟的日志传输,且代码直观,易于维护。

总结: 高级流编程是 Node.js 数据处理能力的精髓。理解如何实现自定义的可读、可写和转换流,并熟练结合压缩、加密等内置流,可以让我们在处理海量数据、构建管道式架构时事半功倍。只要是处理数据流的问题,Node.js 的 Stream 就能提供一种优雅且高性能的解决方案。