Node.js Lambda中实现S3/URL读取流透传至多管道的问题咨询
问题核心原因
- 你定义的
readStream函数没有返回任何值,内部的Promise链未被暴露到外层,调用readStream()拿到的是undefined,无法执行后续.pipe()操作 - 流操作是同步API,你无法等待异步的S3存在性判断完成后再返回流对象,需要提前返回一个
PassThrough透传流作为统一的数据出口 - 原代码中
s3.getObject().createReadStream()未传入Bucket和Key参数,本身就会触发执行错误 - 你测试的
http.get写法错误,回调函数内的return仅作用于回调本身,无法将流返回到外层函数
修复后实现代码
const AWS = require('aws-sdk'); const { PassThrough, pipeline } = require('stream'); const request = require('request'); const sharp = require('sharp'); // 初始化S3实例,Lambda下可直接使用执行角色权限无需额外配置密钥 const s3 = new AWS.S3(); const readStream = ({ Bucket, Key }) => { // 提前返回透传流作为统一出口 const passThrough = new PassThrough(); // 异步判断S3对象是否存在 s3.getObjectMetadata({ Bucket, Key }).promise() .then(() => { // S3对象存在,读取S3流导入透传流 const s3ReadStream = s3.getObject({ Bucket, Key }).createReadStream(); pipeline(s3ReadStream, passThrough, (err) => { if (err) passThrough.destroy(err); }); }) .catch(error => { if (error.statusCode === 404) { // S3对象不存在,读取网页资源流导入透传流 const webReadStream = request.get(`http://example.com/${Key}`); pipeline(webReadStream, passThrough, (err) => { if (err) passThrough.destroy(err); }); return; } // 非404错误直接抛给透传流下游处理 passThrough.destroy(error); }); return passThrough; }; // 调用逻辑和你原有写法一致 readStream({ Bucket: '替换为你的S3桶名', Key: '替换为你的对象键' }) .pipe(sharp().resize(目标宽度, 目标高度).toFormat('png')) .pipe(writeStream); // 替换为你自己的写入流实例
补充说明
- 不需要给
request.get添加await或封装为Promise,request.get本身返回可读流,直接导入透传流即可 - 如果你要替换为原生
http.get,对应分支写法如下:
http.get(`http://example.com/${Key}`, (webReadStream) => { pipeline(webReadStream, passThrough, (err) => err && passThrough.destroy(err)); }).on('error', err => passThrough.destroy(err));
- 这里使用
pipeline而非直接.pipe是为了自动处理错误、资源销毁,避免Lambda出现内存泄漏或异常超时 - 如果你需要透传到多个写入流,直接基于
readStream返回的透传流多次调用pipe即可,透传流会自动将数据广播给所有挂载的写入流
内容的提问来源于stack exchange,提问作者AamirR
相关产品推荐
相关产品推荐

