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

如何利用Node.js Stream填充数组并获取已填充的数组?

解决流处理后数组为空的问题

你的问题核心在于:Node.js的流是异步操作,当你执行完pipe代码后直接访问myArray时,流还没完成所有数据的读取、转换和写入,所以数组是空的。必须等待流的所有操作完成后,才能拿到填充好的数组。

下面给出几种实用的解决方案:

方案1:监听流的finish事件

finish事件会在所有数据都被写入目标流后触发,此时myArray已经填充完成:

const { Readable, Transform, Writable } = require('node:stream')

const myArray = []

// 自定义流类保持不变
class ReadStream extends Readable {
  i = 0
  _read() {
    if (this.i >= 10) {
      this.push(null)
    } else {
      this.i++
      const str = String(this.i)
      const buf = Buffer.from(str, 'ascii')
      this.push(buf)
    }
  }
}

class TransformStream extends Transform {
  _transform(chunk, encoding, callback) {
    const transformed = chunk.toString().toLowerCase()
    callback(null, Buffer.from(transformed))
  }
}

class WriteStream extends Writable {
  _write(chunk, encoding, callback) {
    const write = chunk.toString()
    myArray.push(write)
    callback()
  }
}

// 保存写入流引用,监听finish事件
const writeStream = new WriteStream();
new ReadStream().pipe(new TransformStream()).pipe(writeStream);

// 流完成后再访问数组
writeStream.on('finish', () => {
  console.log('填充完成的数组:', myArray);
  // 在这里执行后续逻辑
});

方案2:用Promise包装流操作(推荐现代写法)

将整个流处理逻辑包装成Promise,用async/await等待结果,代码更清晰:

const { Readable, Transform, Writable } = require('node:stream')

class ReadStream extends Readable {
  i = 0
  _read() {
    if (this.i >= 10) {
      this.push(null)
    } else {
      this.i++
      const str = String(this.i)
      const buf = Buffer.from(str, 'ascii')
      this.push(buf)
    }
  }
}

class TransformStream extends Transform {
  _transform(chunk, encoding, callback) {
    const transformed = chunk.toString().toLowerCase()
    callback(null, Buffer.from(transformed))
  }
}

// 封装成Promise函数
async function getFilledArray() {
  return new Promise((resolve) => {
    const myArray = [];
    const writeStream = new Writable({
      write(chunk, encoding, callback) {
        myArray.push(chunk.toString());
        callback();
      },
      // 所有数据写入完成后触发final
      final(callback) {
        resolve(myArray); // 返回填充好的数组
        callback();
      }
    });

    new ReadStream().pipe(new TransformStream()).pipe(writeStream);
  });
}

// 使用async/await获取结果
(async () => {
  const resultArray = await getFilledArray();
  console.log(resultArray); // 输出 ["1","2",..."10"]
})();

方案3:结合异步迭代器(Generator思路)

利用Node.js流的异步迭代器特性,用for await...of遍历流收集数据,契合你提到的generator使用场景:

const { Readable, Transform } = require('node:stream')

class ReadStream extends Readable {
  i = 0
  _read() {
    if (this.i >= 10) {
      this.push(null)
    } else {
      this.i++
      const str = String(this.i)
      const buf = Buffer.from(str, 'ascii')
      this.push(buf)
    }
  }
}

class TransformStream extends Transform {
  _transform(chunk, encoding, callback) {
    const transformed = chunk.toString().toLowerCase()
    callback(null, Buffer.from(transformed))
  }
}

async function collectWithGenerator() {
  const myArray = [];
  // 将管道后的流转换为异步迭代器
  const stream = new ReadStream().pipe(new TransformStream());

  // 用for await...of遍历(底层基于异步generator)
  for await (const chunk of stream) {
    myArray.push(chunk.toString());
  }

  console.log(myArray);
  return myArray;
}

collectWithGenerator();

内容的提问来源于stack exchange,提问作者Bruno Weber

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 05:45:42