基于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
相关产品推荐
相关产品推荐

