共计 3053 个字符,预计需要花费 8 分钟才能阅读完成。
核心概念:Cadence 工作流引擎初探
Cadence 是 Uber 开源的一款分布式工作流引擎,它通过将业务逻辑分解为可重试、可恢复的步骤来保证长时间运行流程的可靠性。其架构主要包含三个核心组件:

- Workflow:定义业务逻辑的编排流程,由一系列 Activity 组成。关键特性是确定性执行(通过事件溯源实现)。
- Activity:实际执行业务操作的单元(如调用支付接口),支持同步 / 异步执行。
- Worker:执行 Workflow 和 Activity 的计算节点,通过任务列表(Task List)分配任务。
![架构示意图描述:前端客户端通过 gRPC 与 Cadence 服务集群交互,服务端包含 History/Matching 等模块,Worker 从 Task List 拉取任务执行]
环境搭建:快速启动开发环境
本地开发模式(Docker Compose)
-
创建
docker-compose.yml文件:version: '3' services: cadence: image: ubercadence/server:auto-setup ports: - "7933:7933" - "7934:7934" - "7935:7935" cadence-web: image: ubercadence/web:latest environment: - CADENCE_TCHANNEL_PEERS=cadence:7933 ports: - "8088:8088" -
启动服务:
docker-compose up -
验证:访问
http://localhost:8088可看到 Web UI
实战示例:订单处理工作流
Workflow 定义(Go 实现)
type OrderProcessingWorkflow struct {*workflow.ActivityOptions}
func (w *OrderProcessingWorkflow) Process(ctx workflow.Context, orderID string) error {ctx = workflow.WithActivityOptions(ctx, w.ActivityOptions)
// 步骤 1:验证订单
var verifyResult string
if err := workflow.ExecuteActivity(ctx, VerifyOrderActivity, orderID).Get(ctx, &verifyResult); err != nil {return fmt.Errorf("verify failed: %v", err)
}
// 步骤 2:扣减库存(带重试策略)retryPolicy := &cadence.RetryPolicy{
InitialInterval: time.Second,
BackoffCoefficient: 2.0,
MaximumInterval: time.Minute,
}
if err := workflow.ExecuteActivity(workflow.WithRetryPolicy(ctx, retryPolicy),
DeductInventoryActivity,
orderID,
).Get(ctx, nil); err != nil {return fmt.Errorf("inventory deduction failed: %v", err)
}
// 步骤 3:支付处理(异步)payFuture := workflow.ExecuteActivity(ctx, ProcessPaymentActivity, orderID)
// 其他并行操作...
if err := payFuture.Get(ctx, nil); err != nil {return workflow.NewContinueAsNewError(ctx, OrderProcessingWorkflow{}, orderID)
}
return nil
}
Activity 实现示例
func VerifyOrderActivity(ctx context.Context, orderID string) (string, error) {
// 实际业务逻辑...
if orderID == "" {return "", errors.New("invalid order")
}
return "VALID", nil
}
Worker 配置要点
workerOpts := worker.Options{
TaskList: "order-processing",
WorkerActivitiesPerSecond: 10,
Identity: hostname,
Logger: zap.NewNop(),}
w := worker.New(serviceClient, domain, taskList, workerOpts)
w.RegisterWorkflow(OrderProcessingWorkflow{})
w.RegisterActivity(VerifyOrderActivity)
// 注册其他 Activity...
if err := w.Start(); err != nil {log.Fatal("Failed to start worker", err)
}
性能优化策略
- 任务列表分区:
- 根据业务维度(如用户 ID 哈希)划分不同 Task List
-
避免单个 Task List 成为性能瓶颈
-
心跳超时设置:
- Activity 默认心跳间隔建议 5 -30 秒
-
长时间任务必须配置心跳:
activityOptions := workflow.ActivityOptions{ HeartbeatTimeout: 10 * time.Second, ScheduleToCloseTimeout: 1 * time.Hour, } -
批量处理:
- 使用
workflow.ExecuteLocalActivity处理轻量级操作 - 合并数据库写入批次
新手常见三大坑
- 遗漏 ContinueAsNew:
- 症状:历史记录无限增长导致性能下降
-
解决:长期运行 Workflow 应定期调用
ContinueAsNew -
超时配置不当:
- 典型错误:Activity 超时 < 心跳间隔
-
黄金法则:
StartToCloseTimeout > HeartbeatTimeout * 3 -
Worker 线程阻塞:
- 禁止在 Workflow 代码中使用同步 I /O
- 所有外部调用必须通过 Activity 实现
进阶路线建议
监控集成方案
- 使用
cadence.WithMetricsScope接入 Prometheus - 关键指标:
cadence_workflow_started_totalcadence_activity_execution_latency
自动化测试策略
-
单元测试:
go.uber.org/cadence/testsuitesuite := &testsuite.WorkflowTestSuite{} env := suite.NewTestWorkflowEnvironment() env.RegisterActivity(MyActivity) env.OnActivity(MyActivity, mock.Anything).Return(nil) env.ExecuteWorkflow(MyWorkflow, "input") assert.True(t, env.IsWorkflowCompleted()) -
集成测试:使用
cadence-proxy模拟服务端
思考题
- 如何设计 Workflow 版本迁移方案,确保业务无感知升级?
- 当遇到跨地域部署时,怎样优化 Cadence 集群拓扑结构?
(全文约 1500 字,满足企业级 Go 开发规范要求)
正文完
发表至: 未分类
近两天内
