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

基于Node.js实现S3桶间数据同步(含Temporal配置或替代方案)

S3跨桶自动同步方案(Node.js,无AWS付费服务)

一、Temporal 实现方案

1. 本地Temporal集群搭建

无需依赖AWS付费服务,用Docker Compose搭建本地Temporal集群:

  • 创建docker-compose.yml,包含Temporal核心服务(temporal、temporal-ui、cassandra、elasticsearch)
  • 执行docker-compose up -d启动集群,访问http://localhost:8233可查看Temporal UI

2. 核心代码实现

(1)Activity 定义(含重试逻辑)

负责S3文件的下载与上传,配置Temporal的重试策略确保任务完成:

// src/activities.ts
import { S3Client, GetObjectCommand, PutObjectCommand } from "@aws-sdk/client-s3";
import { retry } from "@temporalio/workflow";

const s3Client = new S3Client({ region: "us-east-1" }); // 替换为你的区域

export const syncS3Object = retry(
  async (sourceBucket: string, sourceKey: string, targetBucket: string) => {
    // 拉取源文件
    const sourceRes = await s3Client.send(
      new GetObjectCommand({ Bucket: sourceBucket, Key: sourceKey })
    );
    // 上传至目标桶
    await s3Client.send(
      new PutObjectCommand({
        Bucket: targetBucket,
        Key: sourceKey,
        Body: sourceRes.Body,
        ContentType: sourceRes.ContentType
      })
    );
    console.log(`Synced ${sourceKey} successfully`);
  },
  {
    maximumAttempts: Infinity, // 无限重试直至成功
    backoffCoefficient: 2, // 指数退避
    initialInterval: 1000 // 初始重试间隔1秒
  }
);

(2)Workflow 定义

封装同步任务的执行流程:

// src/workflows.ts
import { proxyActivities } from "@temporalio/workflow";
import type * as activities from "./activities";

const { syncS3Object } = proxyActivities<typeof activities>({
  startToCloseTimeout: "10m" // 单个Activity超时时间
});

export async function s3SyncWorkflow(
  sourceBucket: string,
  sourceKey: string,
  targetBucket: string
) {
  await syncS3Object(sourceBucket, sourceKey, targetBucket);
}

(3)Worker 启动

负责执行Workflow和Activity:

// src/worker.ts
import { Worker } from "@temporalio/worker";
import * as activities from "./activities";

async function runWorker() {
  const worker = await Worker.create({
    workflowsPath: require.resolve("./workflows"),
    activities,
    taskQueue: "s3-sync-queue",
    connection: { address: "localhost:7233" } // 连接本地Temporal集群
  });
  await worker.run();
}

runWorker().catch(err => {
  console.error("Worker crashed:", err);
  process.exit(1);
});

(4)S3事件触发端点

用Express搭建HTTP端点,接收S3的对象创建事件并启动Workflow:

// src/webhook.ts
import express from "express";
import { Connection, WorkflowClient } from "@temporalio/client";
import { s3SyncWorkflow } from "./workflows";

const app = express();
app.use(express.json());
const TARGET_BUCKET = "your-target-bucket"; // 替换为目标桶名

app.post("/s3-webhook", async (req, res) => {
  try {
    const records = req.body.Records;
    for (const record of records) {
      if (record.eventName.startsWith("ObjectCreated:")) {
        const sourceBucket = record.s3.bucket.name;
        const sourceKey = record.s3.object.key;

        // 初始化Temporal客户端
        const connection = await Connection.connect({ address: "localhost:7233" });
        const client = new WorkflowClient({ connection });

        // 启动同步工作流
        await client.start(s3SyncWorkflow, {
          taskQueue: "s3-sync-queue",
          workflowId: `sync-${sourceBucket}-${sourceKey}-${Date.now()}`,
          args: [sourceBucket, sourceKey, TARGET_BUCKET]
        });
      }
    }
    res.status(200).send("Event processed");
  } catch (err) {
    console.error("Failed to trigger workflow:", err);
    res.status(500).send("Internal error");
  }
});

app.listen(3000, () => console.log("Webhook server running on port 3000"));

3. S3事件配置

在源S3桶的控制台中添加事件通知:

  • 事件类型选择“所有创建对象的事件”
  • 目标类型选择“HTTP”,填入你的webhook地址(如http://your-public-ip:3000/s3-webhook)
  • 确保你的服务器能被AWS访问(可使用EC2免费层或ngrok临时暴露本地端口)

二、无Temporal的替代方案

1. 自定义重试 + Webhook

用p-retry库实现重试逻辑,无需依赖Temporal:

// src/sync-service.js
import express from "express";
import { S3Client, GetObjectCommand, PutObjectCommand } from "@aws-sdk/client-s3";
import pRetry from "p-retry";

const app = express();
app.use(express.json());
const s3Client = new S3Client({ region: "us-east-1" });
const TARGET_BUCKET = "your-target-bucket";

async function syncObject(sourceBucket, sourceKey) {
  const sourceRes = await s3Client.send(
    new GetObjectCommand({ Bucket: sourceBucket, Key: sourceKey })
  );
  await s3Client.send(
    new PutObjectCommand({
      Bucket: TARGET_BUCKET,
      Key: sourceKey,
      Body: sourceRes.Body,
      ContentType: sourceRes.ContentType
    })
  );
}

app.post("/s3-webhook", async (req, res) => {
  try {
    const records = req.body.Records;
    for (const record of records) {
      if (record.eventName.startsWith("ObjectCreated:")) {
        const sourceBucket = record.s3.bucket.name;
        const sourceKey = record.s3.object.key;

        // 无限重试直至成功
        await pRetry(() => syncObject(sourceBucket, sourceKey), {
          retries: Infinity,
          factor: 2,
          minTimeout: 1000
        });
      }
    }
    res.status(200).send("OK");
  } catch (err) {
    console.error("Sync failed:", err);
    res.status(500).send("Error");
  }
});

app.listen(3000, () => console.log("Server running on port 3000"));

2. 本地消息队列缓冲(RabbitMQ)

若担心事件丢失,可添加本地RabbitMQ做消息缓冲:

  • Webhook将S3事件推送到RabbitMQ队列
  • 消费者从队列拉取消息,执行同步+重试逻辑
  • 即使服务重启,消息也不会丢失

内容的提问来源于stack exchange,提问作者sooraj

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 12:44:56