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

Node.js Stream无法触发'end'事件导致函数挂起问题排查

问题排查与解决方案

核心原因:未处理的流背压+重复消费冲突

你同时将req管道到两个PassThrough流,但**uploadStream未被消费**,触发Node.js流的背压机制:

  • 当uploadStream内部缓冲区满后,会暂停上游req流的数据推送
  • 这导致etagStream无法收到后续数据,最终processStreamChunks里的for await...of循环因流永远无法结束而挂起

另外,代码同时使用**事件监听(stream.on('data'))和异步迭代器(for await...of)**消费同一个流,两种方式会冲突,干扰流的正常结束逻辑。


修复步骤

1. 确保所有PassThrough流都被消费

如果uploadStream用于AWS SDK上传,必须启动上传逻辑消费该流,示例:

// 确保AWS SDK正在消费uploadStream
const upload = s3.upload({
  Bucket: 'your-bucket',
  Key: 'your-key',
  Body: uploadStream
});

upload.promise().then(() => console.log('上传完成')).catch(err => console.error('上传失败', err));

2. 移除重复的流消费方式

删掉stream.on('data')监听,仅用for await...of消费流,避免冲突:

async function processStreamChunks(stream: Readable) {
  const partSize = 5 * 1024 * 1024;
  let currentPartBuffer = Buffer.alloc(0);
  const partHashes: string[] = [];
  
  // 保留必要的错误和结束监听,移除data监听
  stream.on('end', () => console.log('Stream ended'));
  stream.on('error', (err) => console.error('Stream error:', err));
  stream.on('close', () => console.log('Stream closed'));

  for await (const chunk of stream) {
    currentPartBuffer = Buffer.concat([currentPartBuffer, chunk]);

    while (currentPartBuffer.length >= partSize) {
      const part = currentPartBuffer.slice(0, partSize);
      currentPartBuffer = currentPartBuffer.slice(partSize);
      const partHash = calculateMD5Hash(part);
      partHashes.push(partHash);
    }
  }

  if (currentPartBuffer.length > 0) {
    const partHash = calculateMD5Hash(currentPartBuffer);
    partHashes.push(partHash);
  }

  return partHashes;
}

3. 可选:用pipeline提升流处理健壮性

Node.js的stream.pipeline方法可自动处理背压、错误和流关闭,比手动pipe更可靠:

const { pipeline } = require('stream/promises');

// 并行处理两个流
async function handleStreams(req) {
  const etagStream = new PassThrough();
  const uploadStream = new PassThrough();

  await Promise.all([
    pipeline(req, etagStream, async (source) => {
      const hashes = await processStreamChunks(source);
      console.log('分片哈希计算完成', hashes);
    }),
    pipeline(req, uploadStream, async (source) => {
      await s3.upload({
        Bucket: 'your-bucket',
        Key: 'your-key',
        Body: source
      }).promise();
    })
  ]);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 02:30:30