Cadence技能集合入门指南:从零构建高效工作流引擎

1次阅读
没有评论

共计 3621 个字符,预计需要花费 10 分钟才能阅读完成。

image.webp

为什么选择 Cadence?

Cadence 是 Uber 开源的分布式工作流引擎,它通过持久化函数状态和自动重试机制,解决了传统工作流系统的两大痛点:

Cadence 技能集合入门指南:从零构建高效工作流引擎

  • 长周期任务可靠性 :传统方案如 Airflow 依赖外部数据库记录状态,而 Cadence 内置持久化层,即使进程崩溃也能恢复现场
  • 复杂流程编排 :相比 Temporal 更注重基础架构抽象,Cadence 提供了更丰富的业务级原语(如信号、查询),适合需要深度交互的场景

用一个比喻理解:如果把工作流比作乐高积木,Cadence 不仅提供积木块(基础组件),还附赠组装说明书(状态管理),而其他工具可能只给原材料。

核心概念拆解

1. Workflow(工作流)

持久化的协调逻辑单元,特点包括:

  • 确定性执行 :相同输入永远产生相同输出(不能直接调用随机函数 / 时间函数)
  • 状态持久化 :每次决策后自动保存进度,支持断点续传
  • 超长周期 :理论上可运行数年(通过事件溯源实现)

2. Activity(活动)

实际执行业务逻辑的单元,特点:

  • 非确定性允许 :可以调用外部 API、访问数据库等
  • 自动重试 :内置指数退避等重试策略
  • 超时控制 :支持从秒级到小时级的超时配置

3. Decision Task(决策任务)

Cadence 内部机制,相当于工作流的 ” 大脑 ”:

  1. 调度器将任务放入队列
  2. Worker 领取任务并执行工作流代码
  3. 生成新的决策(如调用 Activity、睡眠等)
  4. 重复上述过程直到工作流完成

Hello World 实战

环境准备

先启动 Cadence 服务端(Docker 版):

docker run -p 7933-7939:7933-7939 -p 7933-7939:7933-7939/udp ubercadence/local-single:0.28.4

工作流定义

// workflow.go
package main

import (
    "go.uber.org/cadence/workflow"
    "go.uber.org/zap"
)

// HelloWorldWorkflow 定义
func HelloWorldWorkflow(ctx workflow.Context, name string) (string, error) {
    // 工作流选项配置
    ao := workflow.ActivityOptions{
        ScheduleToStartTimeout: time.Minute,
        StartToCloseTimeout:    time.Minute,
    }
    ctx = workflow.WithActivityOptions(ctx, ao)

    var result string
    // 执行 Activity
    err := workflow.ExecuteActivity(ctx, SayHello, name).Get(ctx, &result)
    if err != nil {workflow.GetLogger(ctx).Error("Activity failed", zap.Error(err))
        return "", err
    }

    return result, nil
}

Activity 实现

// activity.go
package main

import "context"

// SayHello 实际业务逻辑
func SayHello(ctx context.Context, name string) (string, error) {return "Hello" + name + "!", nil}

Worker 启动

// worker.go
package main

import (
    "go.uber.org/cadence/worker"
    "go.uber.org/zap"
)

func startWorker() {
    // 初始化配置
    workerOptions := worker.Options{Logger: zap.NewNop(),
    }

    // 创建 Worker
    w := worker.New(service.NewClient(service.Options{}),
        "hello-world-domain",
        "hello-world-task-list",
        workerOptions)

    // 注册工作流和 Activity
    w.RegisterWorkflow(HelloWorldWorkflow)
    w.RegisterActivity(SayHello)

    // 启动 Worker
    if err := w.Start(); err != nil {panic("Failed to start worker")
    }
}

客户端调用

// starter.go
package main

import (
    "context"
    "time"

    "go.uber.org/cadence/client"
)

func triggerWorkflow() {
    // 创建客户端
    c := client.NewClient(service.NewClient(service.Options{}), "hello-world-domain", nil)

    // 工作流选项
    options := client.StartWorkflowOptions{
        ID:        "hello-world-workflow",
        TaskList:  "hello-world-task-list",
        ExecutionTimeout: time.Minute,
    }

    // 启动工作流
    we, err := c.StartWorkflow(context.Background(), options, HelloWorldWorkflow, "Cadence")
    if err != nil {panic("Failed to start workflow")
    }

    // 获取结果
    var result string
    err = we.Get(context.Background(), &result)
    if err != nil {panic("Failed to get workflow result")
    }

    println(result) // 输出: Hello Cadence!
}

性能优化三板斧

1. 幂等 Activity 设计

  • 唯一 ID 法 :为每个操作生成唯一 ID,执行前检查是否已处理
    func TransferMoney(ctx context.Context, params TransferParams) error {
        // 基于转账参数生成唯一 ID
        idempotencyKey := fmt.Sprintf(
            "transfer-%s-%s-%d", 
            params.FromAccount, 
            params.ToAccount, 
            params.Amount)
    
        // 检查是否已执行
        if processed, err := checkProcessed(ctx, idempotencyKey); err != nil {return err} else if processed {return nil // 已执行则直接返回}
    
        // 实际业务逻辑...
    }

2. 状态持久化策略

  • 小状态原则 :单个工作流状态不超过 100KB
  • 分离大数据 :将大文件存储在外部存储(如 S3),工作流只记录引用
  • 定期压缩 :对长时间运行的工作流使用 ContinueAsNew 重置状态

3. 分片策略

graph TD
    A[原始任务流] --> B{分片键选择}
    B -->| 用户 ID| C[分片 1: user_1-1000]
    B -->| 地域 | D[分片 2: region_east]
    B -->| 时间 | E[分片 3: Q3_2023]

生产环境避坑指南

错误码处理

错误码 原因 解决方案
EntityNotExistsError 域名未注册 检查 domain.Register API 调用
WorkflowExecutionAlreadyStarted 工作流 ID 重复 使用 UUID 或业务主键作为 ID
DecisionTimeout 决策任务超时 增加 worker 配置的 DecisionTaskTimeout

监控指标

关键 Prometheus 指标:

  • cadence_workflow_started_total
  • cadence_activity_execution_latency
  • cadence_decision_failures_total

推荐告警规则:

- alert: HighWorkflowFailRate
  expr: rate(cadence_workflow_failed_total[5m]) / rate(cadence_workflow_started_total[5m]) > 0.05
  for: 10m

版本升级

采用双阶段升级:
1. 先升级 Worker 保持向后兼容
2. 再升级 Client 使用新特性
3. 最后清理旧版工作流代码

进阶思考

  1. 最终一致性 :如何设计补偿机制处理跨工作流的转账交易?
  2. 动态路由 :能否根据负载自动将工作流转发到不同集群?
  3. 混合云 :怎样实现工作流在公有云和私有环境间的无缝迁移?

通过这个指南,你应该已经掌握了 Cadence 的基础用法。接下来可以尝试用 Cadence 重构公司现有的定时任务系统,体验其自动重试和持久化优势。遇到问题时,记得检查工作流历史记录——这是 Cadence 比其他工具更强大的调试利器。

正文完
 0
评论(没有评论)