Cadence 使用 Skill 入门指南:从零构建可靠的工作流

1次阅读
没有评论

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

image.webp

背景介绍:Cadence 的核心概念

Cadence 是一个由 Uber 开源的分布式工作流编排引擎,旨在解决微服务架构中的长周期业务流程问题。它通过持久化工作流状态和执行历史,确保即使在服务重启或网络中断的情况下,业务流程也能可靠地继续执行。

Cadence 使用 Skill 入门指南:从零构建可靠的工作流

在 Cadence 中,工作流(Workflow)是核心概念,它定义了业务逻辑的执行步骤。每个工作流都有一个唯一的 ID,可以长时间运行(理论上可以是无限期),并且可以包含多个活动(Activity)。活动是工作流中实际执行具体业务逻辑的单元。

痛点分析:新手常见挑战

  1. 状态管理困惑 :新手往往不理解为什么 Cadence 需要显式管理工作流状态,而不是像传统编程那样使用局部变量。

  2. 重试机制不当 :不了解如何正确配置活动重试策略,导致不必要的失败或资源浪费。

  3. 信号处理混乱 :对于如何接收和处理外部信号(Signal)感到困惑,容易造成工作流逻辑混乱。

  4. 超时配置不合理 :没有根据实际业务需求合理设置超时时间,导致工作流过早失败或长时间阻塞。

  5. 版本控制忽视 :忽略工作流版本控制的重要性,在代码更新后遇到兼容性问题。

技术方案:Skill 基本用法与最佳实践

Skill 是 Cadence 提供的客户端库,用于实现工作流和活动。以下是使用 Skill 的基本模式:

  1. 工作流定义 :使用 @workflow.workflow 装饰器定义工作流接口和实现。

  2. 活动定义 :使用 @activity.activity 装饰器定义活动接口和实现。

  3. 工作流客户端 :通过 WorkflowClient 启动和执行工作流。

  4. 工作者配置 :设置 Worker 来监听任务队列并执行工作流和活动。

最佳实践包括:

  • 将工作流和活动分离到不同的模块中
  • 为每个活动定义明确的超时和重试策略
  • 使用合理的任务队列分区策略
  • 实现幂等的活动逻辑
  • 合理使用查询(Query)功能来检查工作流状态

代码示例:简单订单处理工作流

from cadence.activity_method import activity_method
from cadence.worker import Worker
from cadence.workflow import workflow_method, Workflow, WorkflowClient

# 定义活动接口
class OrderActivities:
    @activity_method(task_list="order-task-list", schedule_to_close_timeout_seconds=30)
    async def process_payment(self, order_id: str, amount: float) -> bool:
        """处理支付"""
        pass

    @activity_method(task_list="order-task-list", schedule_to_close_timeout_seconds=60)
    async def send_confirmation(self, order_id: str, email: str) -> bool:
        """发送确认邮件"""
        pass

# 定义工作流接口
class OrderWorkflow:
    @workflow_method(task_list="order-task-list", execution_timeout_seconds=3600)
    async def process_order(self, order_id: str, email: str, amount: float):
        """订单处理工作流"""
        pass

# 实现活动
class OrderActivitiesImpl(OrderActivities):
    async def process_payment(self, order_id: str, amount: float) -> bool:
        # 实际支付处理逻辑
        print(f"Processing payment for order {order_id}, amount: {amount}")
        return True

    async def send_confirmation(self, order_id: str, email: str) -> bool:
        # 实际邮件发送逻辑
        print(f"Sending confirmation for order {order_id} to {email}")
        return True

# 实现工作流
class OrderWorkflowImpl(OrderWorkflow):
    def __init__(self):
        self.activities = Workflow.new_activity_stub(OrderActivities)

    async def process_order(self, order_id: str, email: str, amount: float):
        # 1. 处理支付
        payment_success = await self.activities.process_payment(order_id, amount)
        if not payment_success:
            raise Exception("Payment failed")

        # 2. 发送确认
        await self.activities.send_confirmation(order_id, email)

        return f"Order {order_id} processed successfully"

# 启动工作者
async def start_worker():
    worker = Worker("cadence-system", "order-task-list")
    worker.register_activities_implementation(OrderActivitiesImpl(), "OrderActivities")
    worker.register_workflow_implementation_type(OrderWorkflowImpl)
    await worker.start()

# 启动工作流
async def start_workflow():
    client = WorkflowClient.new_client(domain="sample-domain")
    workflow = client.new_workflow_stub(OrderWorkflow)
    result = await workflow.process_order("order-123", "user@example.com", 99.99)
    print(result)

避坑指南:生产环境常见问题

  1. 活动重试风暴
  2. 问题:活动持续失败导致无限重试
  3. 解决:设置合理的 maximumAttemptsbackoffCoefficient

  4. 工作流历史过大

  5. 问题:长时间运行的工作流产生过多历史记录
  6. 解决:使用 continueAsNew 定期重启工作流

  7. 信号丢失

  8. 问题:工作流未准备好时发送的信号可能丢失
  9. 解决:使用 Workflow.await 确保工作流准备好接收信号

  10. 版本兼容性问题

  11. 问题:更新工作流代码后旧工作流无法继续
  12. 解决:使用 Workflow.getVersion 进行版本控制

  13. 资源泄漏

  14. 问题:工作者未正确关闭导致资源泄漏
  15. 解决:确保在所有情况下调用 worker.shutdown()

总结与思考

通过本文,我们了解了 Cadence 的核心概念,学习了如何使用 Skill 构建可靠的工作流,并探讨了生产环境中可能遇到的问题及其解决方案。

作为进一步探索的方向,可以考虑:

  • 如何设计复杂的工作流决策树?
  • 如何处理跨工作流的协调问题?
  • 如何优化工作流性能,特别是在高并发场景下?

建议读者动手实践以下任务:

  1. 修改示例代码,添加一个取消订单的活动
  2. 尝试实现一个有重试策略的活动
  3. 探索如何在运行时查询工作流状态

Cadence 的强大之处在于它能够将复杂的分布式系统问题简化为工作流编程模型。掌握这些基础后,你将能够构建更加健壮和可靠的分布式系统。

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