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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 17:48:08