如何通过Node.js流实现批量文件下载归档并流式传输至S3
崩溃原因
你当前的实现会在遍历文件列表时瞬间为所有待处理文件创建HTTP下载连接,没有任何并发管控。处理海量文件时,这种写法会快速耗尽系统文件描述符、撑爆内存缓冲区,或是触发源站的连接限流策略,直接导致进程异常退出。node-archiver的append方法接收流参数时,不会自动等待上一个流完全写入归档再处理下一个,Stream自带的背压机制只会在归档写入缓冲区满时暂停写入,不会限制上游同时发起的下载连接数,因此无法解决并发过载的问题。
全流式实现方案
全程基于管道背压机制实现下载、归档、上传的流式处理,不需要额外引入bluebird等依赖,通过轻量的并发槽逻辑即可实现类似Promise.map的并发控制能力,核心逻辑:
- 固定最大并发下载数(建议根据源站限流规则、服务器带宽设置为5-20区间)
- 只有当前一个文件的下载流完全读取、写入归档完成后,才释放并发槽启动新的下载任务
- 全链路添加错误监听,避免单个文件失败导致整个流程崩溃
- 全程不缓存完整文件到内存或本地磁盘,所有数据通过流逐块传输
核心代码
const archiver = require('archiver'); const got = require('got'); const { S3Client } = require('@aws-sdk/client-s3'); const { Upload } = require('@aws-sdk/lib-storage'); // 配置项,根据实际场景调整 const CONCURRENCY_LIMIT = 8; const S3_TARGET_BUCKET = '替换为你的目标桶名'; const S3_TARGET_KEY = '替换为压缩包在S3的存储路径'; const COMPRESS_LEVEL = 5; // 压缩等级1-9,数值越高压缩率越高、CPU消耗越大 // 初始化归档实例 const archive = archiver('zip', { zlib: { level: COMPRESS_LEVEL } }); // 初始化S3上传实例,直接接收归档流作为上传体 const s3Upload = new Upload({ client: new S3Client({ region: '替换为S3所在区域' }), params: { Bucket: S3_TARGET_BUCKET, Key: S3_TARGET_KEY, Body: archive, ContentType: 'application/zip' }, queueSize: 4 // S3分片上传并发数 }); // 全局错误监听,避免未捕获异常导致进程退出 archive.on('error', (err) => { throw new Error(`归档处理异常: ${err.message}`); }); /** * 带并发控制的文件追加逻辑 * @param {Array} fileList 待处理文件列表 */ async function appendFilesWithConcurrencyControl(fileList) { const pendingFiles = [...fileList]; let runningCount = 0; // 任务调度:只要并发槽有空余、还有待处理文件就启动新任务 const runNext = () => { while (runningCount < CONCURRENCY_LIMIT && pendingFiles.length > 0) { runningCount++; const currentFile = pendingFiles.shift(); const { resultFileName, fileUrl } = getFileNameAndUrl(currentFile); // 无有效URL的文件直接跳过 if (!fileUrl) { runningCount--; runNext(); continue; } // 初始化当前文件的下载流 const downloadStream = got.stream(fileUrl, { retry: { limit: 5 }, timeout: { request: 30000 } // 30秒超时,避免僵死连接卡住流程 }); // 单文件下载错误处理:打印日志后跳过该文件,不中断整体流程 downloadStream.on('error', (err) => { console.error(`文件[${resultFileName}]下载失败,已跳过: ${err.message}`); downloadStream.destroy(); runningCount--; runNext(); }); // 等待当前文件完全写入归档后,释放并发槽 archive.append(downloadStream, { name: resultFileName }, (appendErr) => { if (appendErr) { console.error(`文件[${resultFileName}]写入归档失败,已跳过: ${appendErr.message}`); } runningCount--; runNext(); }); } }; // 启动初始任务调度 runNext(); // 等待所有文件处理完成 return new Promise((resolve) => { const checkStatus = setInterval(() => { if (runningCount === 0 && pendingFiles.length === 0) { clearInterval(checkStatus); resolve(); } }, 100); }); } // 主执行流程 async function run() { try { // 启动S3上传监听 const uploadTask = s3Upload.done(); // 按并发限制追加所有文件到归档 await appendFilesWithConcurrencyControl(files); // 所有文件写入完成后,标记归档结束 await archive.finalize(); // 等待S3上传完成 await uploadTask; console.log('压缩包打包并上传S3完成'); } catch (err) { // 异常时销毁流释放资源 archive.destroy(); throw err; } } run();
优化注意事项
- 并发数不要盲目调高:如果源站有单IP访问频率限制,并发过高会触发429拦截;如果服务器出口带宽有限,过高并发不会提升处理速度,反而会增加内存开销
- 可以根据文件大小动态调整并发数:小文件占比高可以适当调大并发,大文件占比高适当降低并发
- 不要把压缩等级设到8以上,高压缩率带来的CPU开销会大幅拖慢整体处理速度,常规场景用5左右的等级性价比最高
- 如果需要支持断点续传、失败重跑,可以在追加文件时记录已处理完成的文件标识,重启脚本时跳过已处理文件即可
内容的提问来源于stack exchange,提问作者Alexander
相关产品推荐
相关产品推荐

