You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何按数量限制动态拆分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}`,
  };
}

关键实现说明

  1. StreamSplitter 自定义流

    • 继承Transform流,工作在对象模式下,通过maxRows配置单文件最大行数
    • 维护当前活跃的PassThrough流,每接收maxRows个对象就关闭当前流并创建新流
    • 触发newStream事件,通知外部处理每个拆分后的流
  2. 多文件写入逻辑

    • 监听newStream事件,为每个新生成的流创建独立的CSV格式化流和文件写入流
    • 使用pipeline处理单一流的管道,确保流的正确关闭和错误捕获
  3. 收尾处理

    • 在_flush阶段关闭最后一个活跃流,确保所有剩余数据都被写入对应文件

内容的提问来源于stack exchange,提问作者Andrei Khotko

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 09:20:14