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

如何在Stream Transform中解析Pending Promise?

解决Transform流中Promise pending的问题

你的问题核心在于transform方法里没有等待异步函数getStuff的Promise完成,就直接推进了chunk,导致chunk.stuff还是pending状态的Promise。下面是修正方案:

关键修改步骤

  • 将transform方法定义为async函数,允许内部使用await等待异步操作完成
  • 用await获取getStuff的实际返回值后,再赋值给chunk.stuff
  • 添加错误捕获,将异步操作的错误传递给stream的回调函数,避免pipeline崩溃
  • 确保所有异步操作完成后再调用callback(),通知stream可以处理下一个chunk

修正后的完整代码

const { pipeline } = require('stream/promises');
const JSONStream = require("JSONStream");
const axios = require('axios');
const fs = require('fs');
const { Transform } = require('stream');

const getStuff = async (id) => {
  const response = await axios.post('/get-stuff', {"id": id});
  return response.data;
}

const addStuff = new Transform({
  writableObjectMode: true,
  async transform(chunk, encoding, callback) {
    try {
      chunk.stuff = await getStuff(chunk.id);
      this.push(JSON.stringify(chunk) + '\n');
      callback();
    } catch (err) {
      callback(err);
    }
  },
})

async function run() {
  await pipeline(
    fs.createReadStream(`./data/input.json`),
    JSONStream.parse('*'),
    addStuff,
    fs.createWriteStream(`./data/output.json`)
  );
  console.log('Pipeline succeeded.');
}

run().catch(console.error);

补充说明

  1. 之前直接赋值chunk.stuff = getStuff(chunk.id)时,getStuff返回的是未完成的Promise,而非实际数据,必须用await等待其resolve。
  2. 将transform设为async函数后,必须在所有异步操作完成后调用callback,否则stream会认为当前chunk已处理完成,导致数据不一致。
  3. 错误处理是必要的:如果getStuff请求失败,将错误传给callback能让pipeline正确捕获错误并终止流程,避免静默失败。
  4. 输出时添加\n是为了让输出的JSON文件每行一个对象,既符合JSONStream的输入格式,也方便后续读取处理。

内容的提问来源于stack exchange,提问作者Sam Luther

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 18:10:29