Cadence Skill 入门指南:从零构建高可靠工作流引擎

1次阅读
没有评论

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

image.webp

为什么需要 Cadence?

在分布式系统中,传统的定时任务方案(如 CronJob)经常面临这些头疼问题:

  • 任务丢失:节点宕机导致任务永远无法执行
  • 重复执行:网络分区可能触发多个实例同时运行
  • 状态管理难:复杂的业务流程需要手动维护执行上下文

Cadence 通过 持久化事件溯源(Event Sourcing)技术,将工作流状态变化记录为不可变事件序列。即使进程崩溃,系统也能从最后记录的事件开始恢复执行。

核心概念三要素

Cadence Skill 入门指南:从零构建高可靠工作流引擎
(示意图说明:Worker 通过轮询 Task Queue 获取任务,执行后更新状态)

  1. Workflow(工作流)
  2. 业务逻辑的协调器,类似状态机
  3. 关键特性:

    • 自动持久化局部变量(无需手动保存)
    • 支持长达数月的长周期执行
  4. Activity(活动)

  5. 实际执行业务操作的单元(如调用支付接口)
  6. 内置:

    • 自动重试机制
    • 心跳检测(Heartbeat)防超时
  7. Task Queue(任务队列)

  8. 解耦 Workflow 和 Activity 的通信通道
  9. 支持权重分配和优先级设置

订单取消实战示例

Java 工作流定义

@WorkflowInterface
public interface OrderWorkflow {
    @WorkflowMethod
    void processOrder(String orderId, Duration timeout);
}

public class OrderWorkflowImpl implements OrderWorkflow {private final ActivityOptions options = new ActivityOptions.Builder()
        .setRetryOptions(new RetryOptions()
            .setInitialInterval(Duration.ofSeconds(1))
            .setMaximumAttempts(3)) // 重试策略
        .build();

    private final OrderActivities activities = Workflow.newActivityStub(OrderActivities.class, options);

    @Override
    public void processOrder(String orderId, Duration timeout) {
        // 启动超时定时器
        CancellationScope cancellationScope = Workflow.newCancellationScope(() -> activities.completeOrder(orderId)
        );

        try {
            // 等待支付完成或超时
            Workflow.await(timeout, () -> activities.isPaymentDone(orderId));
            cancellationScope.cancel(); // 正常支付则取消超时逻辑} catch (TimeoutException e) {activities.cancelOrder(orderId); // 补偿事务
        }
    }
}

CLI 验证步骤

  1. 启动本地 Cadence 服务

    docker run -p 7933:7933 -p 7934:7934 -p 7935:7935 ubercadence/local:latest

  2. 注册工作流类型

    cadence --domain example-domain workflow register --tl order-task-queue --name OrderWorkflow

  3. 触发执行

    cadence --domain example-domain workflow start --tl order-task-queue --workflow_id ORDER_123 --workflow_type OrderWorkflow --execution_timeout 3600

生产环境关键配置

  • 版本升级 :通过workflow.GetVersion() 实现双向兼容

    if (Workflow.getVersion("update1", Workflow.DEFAULT_VERSION, 1) >= 1) {// 新版本逻辑} else {// 旧版本逻辑}

  • 事件压缩:在 cadence-frontend 配置

    historyArchival:
      status: "enabled"
      enableReadFromArchival: true
      retentionPeriod: "720h" # 30 天

  • 监控指标:重点关注

  • cadence_workflow_execution_latency
  • cadence_activity_execution_failed

避坑指南

  1. 幂等性处理:Activity 中必须对写操作做去重

    public void completeOrder(String orderId) {Order order = orderRepository.findById(orderId);
        if (order.getStatus() != Status.NEW) { // 幂等检查
            return;
        }
        // ... 业务逻辑
    }

  2. 上下文传播:在 Async 线程中需手动传递

    try (Scope scope = new ContextPropagationScope()) {CompletableFuture.runAsync(() -> {// 能获取到 Workflow 上下文});
    }

扩展思考

与 Temporal(Cadence 分支版本)的主要差异:

特性 Cadence Temporal
查询 API Workflow.await WorkflowListener
测试框架 需 Mock 内置测试支持
集群管理 需单独部署 集成 K8s Operator

推荐阅读:
Cadence 官方文档
Temporal 迁移指南

写在最后

通过这个订单处理场景的完整实现,相信你已经感受到 Cadence 如何简化分布式协调逻辑。建议从简单的工作流开始,逐步尝试复杂模式如 Saga 事务。遇到问题时,多利用 cadence-cli describe workflow 命令查看内部状态,这对调试非常有帮助。

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