TypeScript中如何在Readable流中替换字符串?
在Node.js流中直接替换字符串并上传S3分块
你从S3拿到的responseBody是个Readable流,要直接在流里替换字符串、输出Buffer用于分块上传,不用把整个文件转成string再处理,核心是用Transform流——这是Node.js stream模块里专门用来做数据中转处理的类型,能边读边处理,完全不用把整个文件加载到内存,适合大文件场景。
核心实现思路
Transform流可以拦截流中的每个数据块(chunk),在里面做字符串替换,同时要注意跨chunk的字符串边界问题:比如要替换的字符串刚好一半在当前chunk末尾,一半在下一个chunk开头,这时候得把残留的部分存起来,和下一个chunk拼接后再处理,避免漏替换。
完整代码示例
const { Transform } = require('stream'); const { S3Client, GetObjectCommand, CreateMultipartUploadCommand, UploadPartCommand, CompleteMultipartUploadCommand } = require('@aws-sdk/client-s3'); // 初始化S3客户端 const s3Client = new S3Client({ region: '你的区域' }); async function processAndUploadS3File(bucket, key, targetBucket, targetKey, searchStr, replaceStr) { // 1. 从S3获取原始流 const s3Response = await s3Client.send( new GetObjectCommand({ Bucket: bucket, Key: key, }) ); const sourceStream = s3Response.Body; // 2. 创建自定义Transform流处理字符串替换 let leftover = ''; // 存储跨chunk的残留字符 const replaceTransform = new Transform({ readableObjectMode: false, writableObjectMode: false, transform(chunk, encoding, callback) { // 把当前chunk转成字符串,加上之前的残留部分 let data = leftover + chunk.toString('utf8'); // 执行替换 data = data.replace(new RegExp(searchStr, 'g'), replaceStr); // 计算残留:如果最后一段可能包含替换字符串的前缀,就留到下一个chunk处理 // 这里简单取替换字符串长度-1的后缀作为残留,你可以根据实际调整逻辑 const leftoverLength = searchStr.length - 1; leftover = data.slice(-leftoverLength); // 把处理后的内容转成Buffer,推出去 this.push(Buffer.from(data.slice(0, data.length - leftoverLength))); callback(); }, flush(callback) { // 处理最后剩下的残留部分 if (leftover) { const finalData = leftover.replace(new RegExp(searchStr, 'g'), replaceStr); this.push(Buffer.from(finalData)); } callback(); } }); // 3. 准备S3分块上传 const multipartUpload = await s3Client.send( new CreateMultipartUploadCommand({ Bucket: targetBucket, Key: targetKey, }) ); const uploadId = multipartUpload.UploadId; const parts = []; let partNumber = 1; const chunkSize = 5 * 1024 * 1024; // 5MB分块,符合S3分块要求(最小5MB,除了最后一块) // 4. 把处理后的流分块上传 return new Promise((resolve, reject) => { let currentChunk = Buffer.alloc(0); replaceTransform.on('data', (data) => { currentChunk = Buffer.concat([currentChunk, data]); // 当当前块达到分块大小,就上传 while (currentChunk.length >= chunkSize) { const partBuffer = currentChunk.slice(0, chunkSize); currentChunk = currentChunk.slice(chunkSize); // 上传当前分块 s3Client.send( new UploadPartCommand({ Bucket: targetBucket, Key: targetKey, UploadId: uploadId, PartNumber: partNumber, Body: partBuffer, }) ).then((partData) => { parts.push({ PartNumber: partNumber, ETag: partData.ETag, }); partNumber++; }).catch(reject); } }); replaceTransform.on('end', async () => { // 上传最后一块剩余数据 if (currentChunk.length > 0) { const partData = await s3Client.send( new UploadPartCommand({ Bucket: targetBucket, Key: targetKey, UploadId: uploadId, PartNumber: partNumber, Body: currentChunk, }) ); parts.push({ PartNumber: partNumber, ETag: partData.ETag, }); } // 完成分块上传 await s3Client.send( new CompleteMultipartUploadCommand({ Bucket: targetBucket, Key: targetKey, UploadId: uploadId, MultipartUpload: { Parts: parts }, }) ); resolve(`文件已成功上传至${targetBucket}/${targetKey}`); }); replaceTransform.on('error', reject); // 把原始流pipe到处理流 sourceStream.pipe(replaceTransform); }); } // 调用示例 // processAndUploadS3File('源桶名', '源文件key', '目标桶名', '目标文件key', '要替换的字符串', '替换后的字符串');
关键细节说明
- 跨chunk处理:用
leftover变量存储每个chunk末尾可能的残留字符,确保不会因为chunk分割漏掉替换。 - Buffer输出:处理后的内容直接转成Buffer,符合S3分块上传的
Body要求,不需要再转成Readable(当然如果需要也可以直接用Transform流作为上传Body,不过分块上传需要拆分Buffer)。 - 分块大小:S3要求分块最小5MB(最后一块除外),代码里设的5MB是合规的,你可以根据需求调整,但不要小于这个值。
内容的提问来源于stack exchange,提问作者MarkBeras
相关产品推荐
相关产品推荐

