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

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事件循环认为没有待处理的异步任务:

  1. 当writeToStream返回等待drain事件的Promise时,Upload实例的上传逻辑如果没有绑定任何监听事件,Node.js不会将其视为活跃的异步任务;
  2. 此时事件循环中只有drain事件的监听,但如果Upload没有占用事件循环资源,Node.js会直接终止进程,不会等待drain触发;
  3. 进程提前退出导致后续写入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 20:30:55