Node.js向S3流式上传时背压未正常排空问题排查
问题描述
使用aws-sdk-js-v3库,通过同一个流将多组数据流式上传至S3的同一个对象,代码如下:
import dotenv from 'dotenv' import {S3} from '@aws-sdk/client-s3' import { Upload } from '@aws-sdk/lib-storage'; import {PassThrough} from 'stream'; import {randomBytes} from 'crypto' import { env } from 'process'; export async function finishS3Stream(upload,stream){ stream.end(); await upload.done(); } export async function writeToStream(stream, data){ // 如果流返回false,表示数据量超过highWaterMark阈值,需要等待drain事件再继续写入 return new Promise((resolve) => { if (!stream.write(data)) { console.log("drain needed") stream.once('drain', resolve); } else { resolve(); } }); } export function createS3Stream(key,bucket) { const client = new S3() const stream = new PassThrough(); stream.on('error', (e) => console.log(e)) const upload = new Upload({ params: { Bucket: bucket, Key: key, Body: stream, }, client, }); return { stream, upload, }; } async function main(){ dotenv.config() const bucket = env.BUCKET const key = env.KEY console.log("creating stream") const {stream,upload} = createS3Stream(key,bucket) const data1 = randomBytes(5242880) console.log('writing data1 to stream') await writeToStream(stream,data1) const data2 = randomBytes(5242880) console.log('writing data2 to stream') await writeToStream(stream,data2) console.log('closing stream') await finishS3Stream(upload,stream) } main()
程序运行后触发highWaterMark阈值,但未等待流排空就以退出码0终止,输出仅显示:
creating stream writing data1 to stream drain needed
需要解决两个问题:
- 如何让程序等待流排空?
- 为何程序未等待Promise解析就提前退出,且未写入下一组数据?
问题原因
程序提前退出的核心原因是Node.js事件循环认为没有待处理的异步任务:
- 当
writeToStream返回等待drain事件的Promise时,Upload实例的上传逻辑如果没有绑定任何监听事件,Node.js不会将其视为活跃的异步任务; - 此时事件循环中只有
drain事件的监听,但如果Upload没有占用事件循环资源,Node.js会直接终止进程,不会等待drain触发; - 进程提前退出导致后续写入
data2、关闭流等逻辑完全没机会执行。
另外,PassThrough默认的highWaterMark仅为64KB(字节模式),远小于你写入的5MB数据块,必然触发drain事件。
解决方法
1. 绑定Upload的进度事件,保持事件循环活跃
给Upload实例绑定httpUploadProgress事件,让Node.js感知到还有正在进行的异步上传任务,不会提前退出:
export function createS3Stream(key,bucket) { const client = new S3() const stream = new PassThrough(); stream.on('error', (e) => console.log(e)) const upload = new Upload({ params: { Bucket: bucket, Key: key, Body: stream, }, client, }); // 绑定进度事件,维持事件循环活跃 upload.on('httpUploadProgress', (progress) => { console.log(`已上传:${progress.loaded}/${progress.total} 字节`); }); return { stream, upload, }; }
2. 优化PassThrough的highWaterMark
设置与数据块大小匹配的highWaterMark,减少drain事件的触发次数:
const stream = new PassThrough({ highWaterMark: 5 * 1024 * 1024 }); // 5MB,和你生成的随机数据大小一致
3. 捕获全局错误,避免静默失败
在main函数末尾添加错误捕获,防止上传过程中出现错误导致进程静默退出:
main().catch(err => { console.error('执行失败:', err); process.exit(1); })
修正后的完整代码
import dotenv from 'dotenv' import {S3} from '@aws-sdk/client-s3' import { Upload } from '@aws-sdk/lib-storage'; import {PassThrough} from 'stream'; import {randomBytes} from 'crypto' import { env } from 'process'; export async function finishS3Stream(upload,stream){ stream.end(); await upload.done(); } export async function writeToStream(stream, data){ return new Promise((resolve) => { if (!stream.write(data)) { console.log("drain needed") stream.once('drain', resolve); } else { resolve(); } }); } export function createS3Stream(key,bucket) { const client = new S3() // 设置匹配数据块大小的highWaterMark const stream = new PassThrough({ highWaterMark: 5 * 1024 * 1024 }); stream.on('error', (e) => console.log(e)) const upload = new Upload({ params: { Bucket: bucket, Key: key, Body: stream, }, client, }); // 绑定进度事件,维持事件循环活跃 upload.on('httpUploadProgress', (progress) => { console.log(`已上传:${progress.loaded}/${progress.total} 字节`); }); return { stream, upload, }; } async function main(){ dotenv.config() const bucket = env.BUCKET const key = env.KEY console.log("creating stream") const {stream,upload} = createS3Stream(key,bucket) // 捕获上传错误 upload.on('error', (err) => { console.error('上传失败:', err); process.exit(1); }); const data1 = randomBytes(5242880) console.log('writing data1 to stream') await writeToStream(stream,data1) const data2 = randomBytes(5242880) console.log('writing data2 to stream') await writeToStream(stream,data2) console.log('closing stream') await finishS3Stream(upload,stream) console.log('上传完成') } main().catch(err => { console.error('执行失败:', err); process.exit(1); })
内容的提问来源于stack exchange,提问作者BigL
相关产品推荐
相关产品推荐

