Node.js流式解析CSV上传S3时触发chunk参数类型非法报错
问题根因
该错误是流传输的数据类型不匹配导致的:
- 你使用的
csv-parse解析流默认运行在对象模式,输出的每个chunk是解析后代表单行数据的普通JS Object,而非原始二进制/字符串数据 - Node.js原生
PassThrough流、AWS S3上传接口要求写入的内容必须是string/Buffer/Uint8Array类型,不支持直接接收JS对象 - 移除
parserStream时链路传输的是S3返回的原始二进制Buffer流,类型匹配因此可以正常运行。
修复方案
核心是在CSV解析流之后追加一个序列化转换流,将解析输出的JS对象重新转换为合法的字符串/二进制格式,再传给下游上传链路,全程保持流式处理不加载全量文件到内存。
方案1:输出CSV格式(和原场景匹配)
使用和csv-parse配套的csv-stringify流做序列化,自动适配表头配置,保持输出为标准CSV格式:
- 先安装依赖:
npm i csv-stringify - 调整流链路代码,插入序列化流,修正后的完整核心逻辑如下:
import { parse } from 'csv-parse'; import { stringify } from 'csv-stringify'; import { PassThrough, Readable } from 'stream'; import { Upload } from '@aws-sdk/lib-storage'; // 省略s3Client、bucketName等已有变量初始化 const objectStream = object?.Body as Readable | undefined; if (objectStream === undefined) { throw new Error('No data'); } const transformationStream = new PassThrough(); // CSV解析流:输出为带修改后表头的行对象 const parserStream = parse({ headers: (headers) => headers.map((header) => header + 'TEST') }) .on('error', (error) => this.log.error(error)) .on('end', (rowCount: number) => this.log.info(`Parsed ${rowCount} rows`)); // CSV序列化流:将行对象转回标准CSV字符串,开启表头输出 const stringifyStream = stringify({ header: true }) .on('error', (error) => this.log.error(error)); // 构建流处理链路:S3下载流 -> CSV解析 -> CSV序列化 -> 透传流 -> S3上传 objectStream .pipe(parserStream) // 如需对每行做自定义转换,可在此处插入自定义Transform流,不要直接监听data事件攒数据破坏背压 .pipe(stringifyStream) .pipe(transformationStream); const upload = new Upload({ client: s3Client, params: { Bucket: this.bucketName, Key: key, Body: transformationStream, }, partSize: 5 * 1024 * 1024, // 配置5MB分片上传,进一步降低大文件内存峰值 }); try { await upload.done(); } catch (error) { this.log.error(error); // 异常时主动销毁所有流,避免资源泄漏 [objectStream, parserStream, stringifyStream, transformationStream].forEach(stream => stream?.destroy()); throw error; }
方案2:输出JSONL格式(无需额外依赖)
如果不需要保留CSV格式,要输出每行一个JSON的JSONL格式,可以自己实现轻量转换流,不需要额外安装依赖:
import { Transform } from 'stream'; // 自定义JSON序列化流,开启对象模式接收JS对象 const jsonlStringifyStream = new Transform({ writableObjectMode: true, transform(row, _, callback) { // 每行序列化为JSON字符串后加换行符 callback(null, `${JSON.stringify(row)}\n`); } }).on('error', (error) => this.log.error(error)); // 替换链路中的stringifyStream即可 objectStream .pipe(parserStream) .pipe(jsonlStringifyStream) .pipe(transformationStream);
注意:不要通过监听
parserStream的data事件手动拼接全量数据,会破坏Node.js流的背压机制,大文件场景下依然会出现内存溢出问题。
内容的提问来源于stack exchange,提问作者db9
相关产品推荐
相关产品推荐

