微服务编排
协调并发服务调用,接收执行中的状态更新,并通过外部事件或超时兜底完成流程。
产品需求
流程以一个数据值启动,并且必须:
- 调用 API 1,然后把输入持久化到 Data Attribute。
- 并发启动 CallAPI2 和 CallAPI3。
- 让 CallAPI2 读取当前数据、调用 API 2,并结束自己的分支。
- 让 CallAPI3 等待一条 Ready Channel 消息或一个 24 小时的持久 Timer,任一条件满足即可继续。
- 允许调用方在 Flow 运行期间调用 Swap RPC。RPC 必须原子地返回旧值并写入新值。
- 收到 Ready 消息后,使用最新数据调用 API 3,然后完成 Flow。
- 如果 Timer 先触发,则调用 API 3,继续到 CallAPI4,使用最新数据调用 API 4,然后完成。
这段等待是持久的。Worker 或 Dex 重启不会丢失 Channel 消息、Timer 截止时间、持久化数据或当前 Step。Ready 消息通常由另一个服务在异步工作完成后发送。如果消息始终没有到达,Timer 会选择兜底路径。
潜在故障包括:
- 下游 API 可能不可用、超时或拒绝输入。Step 必须上报故障,让配置的重试和失败策略处理它。
- 下游调用可能已经成功,但响应丢失。Dex 可能重试 Step,因此生产环境中的 API 操作必须具备幂等性。
- Swap RPC 或 Ready 发布可能指向不存在或已经完成的 Flow,并返回客户端错误。调用方应保留 Flow ID 并处理该响应。
- Ready 消息可能一直没有到达。这属于预期结果而非进度丢失;持久 Timer 最终会让 Flow 进入 CallAPI4。
- 如果 Swap 恰好发生在两个分支读取 Attribute 之间,CallAPI2 和 CallAPI3 可能看到不同的值。流程保证每次读取都是持久的,但不保证并发分支共享同一个快照。
示例使用本地依赖 stub。生产适配器还应为每个 API 配置超时、幂等键,以及适当的重试或恢复策略。
Flow 设计
这个交互式 definition graph 由可运行的 Python Flow 生成。
这个 Flow 包含一个 fan-out 点和两种完成路径:
CallAPI1 ─┬─▶ CallAPI2 ─▶ DeadEnd
│
└─▶ Wait(Ready | 24h Timer) ─▶ CallAPI3 ─┬─▶ Complete
└─▶ CallAPI4 ─▶ Complete
(仅 Timer)
它组合了以下设计模式:
- 静态并行 Step 并发启动两个已知的服务调用,不把两个操作放进同一个重试边界。
- 持久外部事件门控通过 Wait.anyOf 组合 Channel 和 Timer。Channel 处理正常回调,Timer 保证等待时间有上限。
- 超时兜底只把超时的执行送到 CallAPI4。正常的 Channel 路径会跳过这个兜底步骤。
- 请求—响应交互把 Swap 暴露为强类型 RPC,让调用方在一次操作中更新持久状态并收到旧值。
- 持久共享状态把当前 payload 存入 Attribute。两个并发分支在执行时读取数据,而不是在图中一直携带可能过期的副本。
- 分支终止让 CallAPI2 返回 DeadEnd,由等待分支提供业务完成结果。
每个下游 API 都有独立 Step。因此重试 API 3 不会重复 API 1 或 API 2,Step 执行图也能准确显示哪个调用正在等待、重试或已经完成。
核心实现
下面的标签页展示 Flow 签名、拓扑、持久化 schema、Swap RPC、并发 fan-out,以及 Channel 或 Timer 等待。每个链接还指向包含依赖 stub、HTTP controller 和端到端集成测试的完整示例。
class OrchestrationFlow(Flow[str]):
data = Attribute("data", str)
ready = Channel[None]("Ready", type(None))
def __init__(self, service: MyDependencyService) -> None:
self.service = service
self.call_api4 = CallAPI4(service, self.data)
self.call_api3 = CallAPI3(service, self.data, self.ready, self.call_api4)
self.call_api2 = CallAPI2(service, self.data)
self.call_api1 = CallAPI1(service, self.data, self.call_api2, self.call_api3)
def get_steps(self) -> StepList[str]:
return StepList.start_step(self.call_api1).other_steps(
self.call_api2,
self.call_api3,
self.call_api4,
)
def get_persistence_schema(self) -> PersistenceSchema:
return PersistenceSchema.of(self.data, self.ready)
@rpc
def swap(self, context: Context, new_data: str) -> RPCResult[str]:
old_data = self.data.get(context)
self.data.set(context, new_data)
return RPCResult(old_data)
class CallAPI1(Step[str]):
def execute(self, context: Context, input: str) -> StepDecision:
self.service.call_api1(input)
self.data.set(context, input)
return go_to_many(
StepMovement.of(CallAPI2, None),
StepMovement.of(CallAPI3, None),
)
class CallAPI3(Step[None]):
def wait_for(self, context: Context, input: None) -> Wait:
return Wait.any_of(
Timer.by_duration(timedelta(hours=24)),
self.ready.for_one(),
)
def execute(self, context: Context, input: None) -> StepDecision:
value = self.data.get(context)
self.service.call_api3(value)
if context.has_timer_fired():
return go_to(CallAPI4, None)
return graceful_complete(value)
例子: examples/python/dex_examples/products/microservices/orchestration_flow.py
演示
启动 Dex 和任意语言的示例服务,然后通过 examples playground 启动 Microservice orchestration。在下面这次运行中,Flow 以 test initial data 启动,Swap RPC 把它替换为 updated in-flight data,随后一条 Ready 消息结束等待。CallAPI1、CallAPI2 和 CallAPI3 均已完成;因为 Channel 在 24 小时 Timer 之前满足条件,所以 CallAPI4 没有执行。
