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

Node.js中S3流复制后调用批量写入致records为空问题

问题根源分析

你的问题出在流的背压机制和多消费者的冲突上:

  • 你通过duplicateStream把原始S3流同时pipe到两个PassThrough流,这两个流共享原始流的数据推送。
  • processRecordCount函数在读取第一行数据后,执行了stream.unpipe(parser)和parser.end()——这里的stream是其中一个PassThrough实例(streamForRecordCount)。取消流与parser的连接后,该PassThrough的输出端失去消费者,根据Node.js流的背压规则,它会向原始S3流发送暂停推送的信号。
  • 原始S3流只要有一个pipe目标处于暂停状态,就会停止向所有目标推送数据。这直接导致另一个PassThrough流(streamForEntireRecords)无法接收到后续CSV行,最终processEntireRecords返回空数组。
  • 移除processRecordCount时,只有streamForEntireRecords一个消费者,原始流会正常推送所有数据,因此records能正常生成,batchWriteItems也能正常工作。
解决方案

方案一:修改processRecordCount,避免提前终止流连接

不需要取消pipe,仅在拿到第一行后忽略后续数据即可。这样streamForRecordCount会继续接收原始流的所有数据(只是我们不处理),不会触发背压暂停:

async function processRecordCount(stream) {
  return new Promise((resolve, reject) => {
    let firstLineProcessed = false;
    const results = [];

    const parser = csvParser();

    stream.on('error', (err) => {
      console.error('Error with the input stream:', err);
      reject(err);
    });

    parser.on('error', (err) => {
      console.error('Error parsing CSV:', err);
      reject(err);
    });

    parser.on('end', () => {
      if (results.length === 0) {
        reject(new Error('CSV is empty'));
      } else {
        resolve(results[0]); 
      }
    });

    stream.pipe(parser);

    parser.on('data', (row) => {
      if (!firstLineProcessed) {
        results.push(row);
        firstLineProcessed = true;
        // 移除 stream.unpipe(parser) 和 parser.end(),仅忽略后续数据
      }
    });
  });
}

方案二:先缓存原始流数据,再生成独立流处理

如果CSV文件大小不大,可以先把S3流转换成Buffer,再基于Buffer生成两个独立的可读流,让两个流完全互不影响:

// 实现流转Buffer的工具函数
async function streamToBuffer(stream) {
  return new Promise((resolve, reject) => {
    const chunks = [];
    stream.on('data', chunk => chunks.push(chunk));
    stream.on('end', () => resolve(Buffer.concat(chunks)));
    stream.on('error', reject);
  });
}

// 主流程修改
const command = new GetObjectCommand(getObjectParams);                                                                                                                                                                              
const data = await s3Client.send(command);                                                                                                                                                                                          
const buffer = await streamToBuffer(data.Body);
// 生成两个独立的可读流
const { Readable } = require('stream');
const streamForRecordCount = Readable.from(buffer);
const streamForEntireRecords = Readable.from(buffer);

const recordedCount = await processRecordCount(streamForRecordCount);                                                                                                                                                               
const records = await processEntireRecords(streamForEntireRecords);     
await batchWriteItems(records, tableName);
console.log("records", records);

方案三:并行处理两个流

不要串行await两个处理函数,而是同时启动两个流的处理,避免其中一个流的提前结束影响另一个:

const [streamForRecordCount, streamForEntireRecords] = duplicateStream(stream);
// 并行启动两个处理任务,用Promise.all等待结果
const [recordedCount, records] = await Promise.all([
  processRecordCount(streamForRecordCount),
  processEntireRecords(streamForEntireRecords)
]);
await batchWriteItems(records, tableName);
console.log("records", records);
验证说明

修改后,streamForEntireRecords能正常接收所有CSV数据,processEntireRecords返回完整的records数组,batchWriteItems也能正常写入数据,最终console.log会输出预期内容。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 03:07:34