跳到主要内容

微服务编排

协调并发服务调用,接收执行中的状态更新,并通过外部事件或超时兜底完成流程。

产品需求

流程以一个数据值启动,并且必须:

  1. 调用 API 1,然后把输入持久化到 Data Attribute。
  2. 并发启动 CallAPI2CallAPI3
  3. CallAPI2 读取当前数据、调用 API 2,并结束自己的分支。
  4. CallAPI3 等待一条 Ready Channel 消息或一个 24 小时的持久 Timer,任一条件满足即可继续。
  5. 允许调用方在 Flow 运行期间调用 Swap RPC。RPC 必须原子地返回旧值并写入新值。
  6. 收到 Ready 消息后,使用最新数据调用 API 3,然后完成 Flow。
  7. 如果 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 之间,CallAPI2CallAPI3 可能看到不同的值。流程保证每次读取都是持久的,但不保证并发分支共享同一个快照。

示例使用本地依赖 stub。生产适配器还应为每个 API 配置超时、幂等键,以及适当的重试或恢复策略。

Flow 设计

这个交互式 definition graph 由可运行的 Python Flow 生成。

Definition graph

OrchestrationFlow

Valid

python · examples/python/dex_examples/products/microservices/orchestration_flow.py

这个 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 消息结束等待。CallAPI1CallAPI2CallAPI3 均已完成;因为 Channel 在 24 小时 Timer 之前满足条件,所以 CallAPI4 没有执行。

已完成的 Microservice orchestration Step 执行图

相关内容