• 后端
  • 任务调度
  • 工作流自动化
  • 微服务

【免费下载链接】cadence

Cadence is a distributed, scalable, durable, and highly available orchestration engine to execute asynchronous long-running business logic in a scalable and resilient way.

项目地址: https://gitcode.com/gh_mirrors/cad/cadence
点击查看 免费下载

导读

Workflow Shadowing(工作流阴影测试)是 Cadence 为解决工作流定义非向后兼容变更(non-backward compatible change)所引入的一种全新 Worker 模式:它在不启动正常 activity/decision worker 组件的前提下,专门拉取生产环境的工作流历史并基于新版本的工作流定义进行回放(replay),从而在本地开发、staging/pre-prod/canary 等预生产阶段提前暴露非确定性错误(non-deterministic error),显著降低线上事故的发生率。本文以 docs/design/workflow-shadowing/2547-workflow-shadowing.md 为骨架,结合当前 Cadence 仓库中的前端 API 定义等源码证据,完整拆解该设计的背景、方案、备选路线、Worker 配置项与指标体系,帮助读者理解如何在发布流水线中落地这一兼容性校验机制。

背景:非确定性错误从何而来

Cadence 客户端库(Golang / Java)在构建当前工作流状态时,依赖两份关键输入:

  • 工作流历史(workflow history):由 Cadence server 持久化,且不可变;
  • 工作流定义(workflow definition):由用户通过 Worker 部署控制,会随版本迭代而变化。

当客户端使用某个版本的工作流历史、却用另一版本的工作流定义来重建状态时,就可能发生非确定性错误。原设计文档给出了最典型的场景:

版本工作流定义工作流历史
V11. 启动工作流
2. 启动 Activity A
3. 完成工作流
1. 工作流已启动
2. Activity A 已启动
V21. 启动工作流
2. 启动 Activity B
3. 完成工作流
1. 工作流已启动
2. Activity B 已启动

当客户端用 V1 的历史 + V2 的定义重建状态时,它看到历史中存在"Activity A 已启动"事件,于是期望在状态中找到 Activity A;但 V2 定义里根本没有 Activity A(定义中指定的是 Activity B)。双方对不上,就会在 workflow(decision)task 处理过程中抛出非确定性错误。

依据非确定性错误处理策略的不同,工作流要么立刻失败,要么被阻塞直到人工介入。原文档明确指出:这类错误曾在多个客户环境中引发事故,而更麻烦的是理解与缓解这类错误通常非常耗时——盲目回滚糟糕的部署只会让情况更糟,因为新定义产生的新历史同样与旧定义不兼容。正确的缓解手段既需要人工操作,又需要深入理解错误成因,这大大拉长了故障恢复时间。

目标与非目标

  • 目标(Goals):在预生产环境中检测非确定性错误、通知客户,从而减少由非确定性错误引发的事故数量。
  • 非目标(Non-Goals):
    • 不追求 100% 杜绝生产环境的非确定性错误(该方案并非"防弹"方案);
    • 不提供针对既有非确定性错误的缓解/修复方案(本设计只聚焦检测)。

方案提案:Shadowing 模式的 Replay Worker

设计的核心提案是:为 Worker 新增一种名为 shadowing 的运行模式,复用 Golang 与 Java 客户端中已有的工作流回放测试框架(replay test framework),仅用"新版本的工作流定义 + 生产环境的工作流历史"执行回放测试。

三步执行流程

当用户 Worker 以 shadowing 模式运行时:

  1. 普通 activity / decision worker 组件不再启动,取而代之启动一个专门的 replay worker;
  2. replay worker 调用 ScanWorkflowExecutions API 获取一批工作流可见性记录(visibility records);
  3. 对每条记录,通过 GetWorkflowExecutionHistory API 从 Cadence server 拉取对应工作流的历史,再用客户端侧的回放测试框架,将历史与新版本的工作流定义进行回放比对;回放失败时,发出 metrics 与日志消息,通知用户"新定义中存在非确定性变更"。

上述两个 API 在当前仓库的 client/frontend/interface.go 中均有对应定义:GetWorkflowExecutionHistory(context.Context, *types.GetWorkflowExecutionHistoryRequest, ...yarpc.CallOption) (*types.GetWorkflowExecutionHistoryResponse, error)(第 50 行)与 ScanWorkflowExecutions(context.Context, *types.ListWorkflowExecutionsRequest, ...yarpc.CallOption) (*types.ListWorkflowExecutionsResponse, error)(第 78 行),可见该流程走的是标准的 frontend 服务接口。

运行位置的灵活性

只要 Worker 能与 Cadence server 通信以获取可见性记录与工作流历史,就能以 shadowing 模式运行。这意味着用户可以选择:

  • 本地环境运行:配合快速开发迭代;
  • 作为部署流程的一环:在 staging/pre-prod 环境的专用主机上运行,以获得更好的覆盖并捕获罕见场景。

用户还可以自定义"要阴影哪些工作流"。例如:不使用 query 功能时,可以跳过对 closed 工作流的阴影;只改动了某个工作流类型时,只需阴影并检查该类型的工作流历史。具体选项见下文"Shadow Worker 选项"一节。

部署流水线全景

Cadence workflow shadowing 部署流水线:本地集成测试 → staging/pre-prod shadowing worker → 生产 worker,三个阶段均对接生产 Cadence server

上图(docs/design/workflow-shadowing/2547-deployment-pipeline.png)展示了完整的开发流水线:本地环境的 Local integration tests(Local shadowing worker) 与预发布环境的 Shadowing worker 都直接对接 Production Cadence server 拉取真实历史做回放校验,验证通过后再进入生产环境的 Prod worker 部署,从而实现"先阴影、后上线"的发布节奏。

方案的主要优点

  • 简单(Simplicity):Cadence server 端几乎无需改动;
  • 多环境可用:本地、staging/pre-prod/canary 环境均可运行;
  • 通用性强:适配内部部署流程、CI/CD pipeline,也适配开源用户。

方案的已知代价

  • 用户侧需要一定投入:需要搭建发布流水线与阴影环境,Cadence 需为此提供使用指南;
  • server 与 ElasticSearch 集群负载上升:新增的 GetWorkflowExecutionHistory 与 ScanWorkflowExecutions 调用会带来额外负载,但可通过服务端限流器(rate limiter)控制(原文档在 "Open Questions" 一节中给出了负载估算思路);
  • 专用 shadow worker 在回放测试结束后会闲置。

其他被评估过的备选方案

设计文档还系统评估了三条替代路线,并解释了为何未采用:

方案一:基于 task mocking 的历史生成(History Generation with task mocking)

思路与本地回放测试框架类似,但不再要求用户从生产拉取历史,而是让用户定义可能的工作流输入、activity 输入/输出、signal 等,基于某个工作流版本按不同取值组合自动生成多条工作流历史,并注入随机错误覆盖失败场景;代码变更后,用这些生成的历史对新定义做回放比对。

优点:可视为"测试用例自动生成"的自动化测试,能在本地完成校验,开发迭代更快,可作为 vNext 阶段改进本地开发体验的方向。

缺点:

  • 用户需学习一套新的 API 来定义各类输入,新增 activity 或 signal channel 时还要同步更新测试;
  • 覆盖率存疑:许多输入来自上游服务(如 Kafka),难以穷举所有参数组合;
  • 实现上是语言相关的:需要为不同客户端库分别重新设计 API 与实现。

方案二:基于归档(部分)历史的集成测试(Integration test with archived history)

复用现有历史归档(archival)代码路径,把符合某查询条件的工作流历史 dump 到 blob storage;用户改完代码后,连接 blob storage 在本地跑集成测试回放全部归档历史。由于部分 blob storage 在条目量大时不支持 List 操作,可能还需依赖 visibility archival 获取 workflowID/runID 列表。

优点:server 端改动少(历史/可见性归档代码路径可复用);避免 worker 闲置问题;不增加 Cadence server 负载(测试框架直连存储)。

缺点:

  • 需要确定 blob storage 的归属与管理方:
    • 若 Cadence 管理:部分历史含 PII 数据,GDPR 约束下不能归档过久,还需一套存储管理 API、注册用户回调;
    • 若用户管理:会增加用户负担(如按 GDPR 要求定期清理、调整查询只回放近期工作流以保证测试时长合理);
  • 需要额外搭建历史归档与可见性归档两套系统;
  • 本地跑全量归档回放耗时过长,不适合本地运行;
  • 测试周期长且不在部署流程内,用户可能跳过测试,直接带病上线;
  • 开源视角下用户必须先具备归档系统(当时仅支持 S3 与 GCP),否则要自行实现 archiver;客户端侧虽可提供 blob storage 接口,但用户仍需自行实现从存储拉取历史/可见性记录的逻辑;
  • 将 blob storage 作为外部依赖后,依赖不可用则方案失效。

方案三:基于复制的回放测试(Replication based replay test)

复用复制(replication)栈,把 decision task 分发到 standby 集群的 worker 上执行回放校验。该方案仅对全局域(global domain)用户有效,覆盖面不足,因此被排除。

实现细节

Worker 实现的两条路径

shadowing 的三步流量(获取可见性记录 → 获取历史 → 回放测试)既可以用原生代码实现,也可以用 Cadence 工作流实现:

  • 本地开发环境:提供一套 shadowing 测试框架,让用户更早验证变更、加速开发;这类集成测试也容易接入现有发布流水线,作为代码落地时的强制检查。
  • staging/生产环境:倾向用 Cadence 工作流来控制 shadowing,理由有四:
    1. 可扩展性(Scalability):支撑大规模 shadowing 流量;
    2. 可扩展性(Extensibility):容易扩展到不同语言;
    3. 可维护性(Maintenance):无需为 shadowing 模式维护专门的 worker 代码;
    4. 可见性(Visibility):shadowing 进度与结果可直接通过 Cadence Web UI 与指标仪表盘查看。

shadowing 处理过程中会持续发出 metrics,因此可以接入发布流水线,在部署到生产前检查工作流兼容性。代价是:用户需要启动一个 Cadence server 来运行 shadow 工作流。

Shadow Worker 选项

以下选项将同时加入 Golang 与 Java 客户端的 Worker Option 结构体:

  • EnableShadowWorker:布尔值,表示 worker 是否以 shadow 模式运行;
  • ShadowMode:
    • Normal:扫描完全部工作流后停止 shadowing(默认值);worker 重启或重新部署时触发新一轮迭代;
    • Continuous:持续 shadowing,直到满足退出条件。因为 open 工作流会持续进展,两次回放可能得到不同结果;该模式适合工作流会长时间阻塞在 timer 上的用户;
  • Domain:用户域(domain)名称;
  • TaskList:shadowing activity worker 轮询的 task list 名称;
  • ShadowWorkflowQuery:需要回放的所有工作流的可见性查询语句;若指定,则忽略下列三个选项;
  • ShadowWorkflowTypes:需要检查的工作流类型列表;
  • ShadowWorkflowStartTimeFilter:工作流开始时间的时间范围;
  • ShadowWorkflowStatus:工作流关闭状态列表,只有关闭状态在列表中的工作流才会被检查;同时提供 open 工作流选项,并作为默认值;
  • SamplingRate:浮点数,表示对 scan 返回的工作流执行回放的百分比;
  • Concurrency:回放检查的并发度,即 shadow 工作流中并行回放 activity 的数量。

本地 Shadow 测试选项

本地测试框架提供的选项与 Shadow Worker 选项类似,但测试中扫描与回放只有单线程,因此仅接受 concurrency = 1;同时移除 EnableShadowWorker、Domain、TaskList 三个选项。

工作流代码归属:为什么放在 server 侧

设计决定将 shadowing 工作流代码保留在 Cadence server 侧,理由如下:

  1. Java 与 Go 客户端的 shadowing 逻辑完全一致,放在 server 侧只需实现一次,无需维护两份工作流定义并保持同步;
  2. 既然 server 拥有工作流代码,Cadence 团队可随时更新定义、甚至做非向后兼容的修改,而无需要求客户升级客户端依赖。

相应代价是:

  1. Cadence 团队需要提供 worker 来处理 shadowing 工作流产生的 decision task;
  2. 所有客户的 shadowing 工作流会共处一个 domain,由于 domain 总资源有限,可能互相影响执行。

不过,考虑到 shadowing 工作流产生的负载很小、且不太可能大量用户同时运行,原文档判断这些缺点不足为虑。

指标(Metrics)

shadow 工作流将发出以下指标:

  1. 成功 / 跳过 / 失败的工作流回放数量;
  2. 每次回放测试的延迟;
  3. shadow 工作流的 Start / Complete / Continue-as-new 计数;
  4. 每个 shadow 工作流的延迟。

这些指标配合 Cadence 指标体系(可参考 common/metrics/defs.go 中的指标定义方式)即可接入发布流水线做自动化的兼容性闸门。

源码现状说明

需要说明的是,当前仓库中暂未包含 shadowing worker 的最终实现代码——本设计文档(原 issue 引用为 #2547)描述的是该功能的方案提案:按照设计,shadowing 工作流代码归属 server 侧、回放框架归属各客户端库,实际落地分散在 Cadence server 与 Golang/Java 客户端仓库中。读者可从本文引用的 frontend 接口定义 与部署流水线示意图(2547-deployment-pipeline.png)确认核心数据通路的设计。

总结

Cadence Workflow Shadowing 以极小的 server 端改动,把"回放测试"从开发者的本地工具升级为发布流水线中的强制兼容性检查:通过 ScanWorkflowExecutions + GetWorkflowExecutionHistory 拉取生产历史、复用客户端回放框架、以 Cadence 工作流驱动并发回放,并在预生产阶段提前暴露非确定性错误。配合 Normal/Continuous 两种 shadow 模式、可见性查询过滤、采样率与并发度控制,以及成功/失败/延迟等指标,用户可以在不改动生产定义的前提下,用最贴近线上真实数据的方式守护工作流定义的向后兼容性。

  • 后端
  • 任务调度
  • 工作流自动化
  • 微服务

【免费下载链接】cadence

Cadence is a distributed, scalable, durable, and highly available orchestration engine to execute asynchronous long-running business logic in a scalable and resilient way.

项目地址: https://gitcode.com/gh_mirrors/cad/cadence
点击查看 免费下载
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐