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
相关产品推荐
相关产品推荐

