如何按数量限制动态拆分Node.js Readable流为多个流
实现按指定行数拆分对象流并生成多个CSV文件
现有输出对象的Readable流,需按指定行数n拆分,每个拆分块对应一个CSV文件(通过@fast-csv/format处理)。例如n=10、总对象数32时,生成3个含10行、1个含2行的CSV文件。以下是可配置n的自定义流实现方案:
完整实现代码
const { Readable, Transform, PassThrough } = require('node:stream'); const { format } = require('@fast-csv/format'); const fs = require('fs'); const path = require('path'); const { pipeline } = require('node:stream/promises'); // 自定义分流器:按指定行数拆分对象流,生成多个可读流 class StreamSplitter extends Transform { constructor(options) { super({ objectMode: true, ...options }); this.maxRows = options.maxRows; this.currentRowCount = 0; this.currentStream = null; this.streamIndex = 0; this.streams = []; // 存储所有生成的可读流 // 初始化第一个流 this._createNewStream(); } _createNewStream() { this.currentStream = new PassThrough({ objectMode: true }); this.streamIndex += 1; this.streams.push(this.currentStream); // 触发新流创建事件,供外部处理 this.emit('newStream', this.streamIndex, this.currentStream); } async _transform(chunk, encoding, callback) { this.currentStream.write(chunk); this.currentRowCount += 1; // 达到最大行数时,关闭当前流并创建新流 if (this.currentRowCount >= this.maxRows) { this.currentStream.end(); this.currentRowCount = 0; this._createNewStream(); } callback(); } _flush(callback) { // 关闭最后一个活跃流,确保剩余数据写入 if (this.currentStream && !this.currentStream.destroyed) { this.currentStream.end(); } callback(); } } (async () => { try { const maxRowsPerFile = 10; const readable = new Readable({ objectMode: true }); const streamSplitter = new StreamSplitter({ maxRows: maxRowsPerFile }); // 监听新流创建事件,每个流对应一个CSV文件 streamSplitter.on('newStream', async (index, stream) => { const csvStream = format({ headers: true }); const fileStream = fs.createWriteStream( path.join(process.cwd(), `test_${index}.csv`), { flags: 'w' } ); await pipeline(stream, csvStream, fileStream); console.log(`文件 test_${index}.csv 写入完成`); }); // 推送测试数据 const objects = createObjects(32); objects.forEach(obj => readable.push(obj)); readable.push(null); // 将原始流接入分流器 await pipeline(readable, streamSplitter); console.log('所有文件处理完成!'); } catch (err) { console.error(err); } })(); // 工具函数:生成测试对象 function createObjects(n) { const objects = []; for (let i = 0; i < n; i++) { objects.push(createObject(i)); } return objects; } function createObject(i) { return { id: i, name: `Obj #${i}`, }; }
关键实现说明
StreamSplitter 自定义流
- 继承
Transform流,工作在对象模式下,通过maxRows配置单文件最大行数 - 维护当前活跃的
PassThrough流,每接收maxRows个对象就关闭当前流并创建新流 - 触发
newStream事件,通知外部处理每个拆分后的流
- 继承
多文件写入逻辑
- 监听
newStream事件,为每个新生成的流创建独立的CSV格式化流和文件写入流 - 使用
pipeline处理单一流的管道,确保流的正确关闭和错误捕获
- 监听
收尾处理
- 在
_flush阶段关闭最后一个活跃流,确保所有剩余数据都被写入对应文件
- 在
内容的提问来源于stack exchange,提问作者Andrei Khotko
相关产品推荐
相关产品推荐

