跳到主要内容

等待 Step 完成

当一个 Flow 会在后台工作结束前产生 durable、可供客户端使用的里程碑时,使用 WaitForStepCompletion。调用方启动 Flow、等待该 Step,然后读取需要的结果。后续 Step 可以继续执行,而无需一直占用调用方。

WaitForStepCompletionFlow 示例会在 PersistData 中持久化记录,再启动 BackgroundWork。它的 controller 在查询已持久化记录前等待 PersistData。只有当一个 Step 已提交了使响应有效的状态后,才应把它选作这个边界。

Dex 会从 Step execution 派生 Request ID,因此等待同一里程碑的调用方会共享一个已接受的 durable Update。普通的无限 handler 等待应让 MaximumWaitTime 保持为 0。本地 HTTP 或 client deadline 仍可限制一次响应的时长,而不会取消该已接受的 handler。只有当一个 Flow 可能积累被放弃或很少完成的 Step 等待时,才设置正值。过期会释放 in-flight slot,但继续等待会创建新的 -N Update generation。容量和成本的取舍请参阅选择 durable handler 的生命周期

核心实现

@blueprint.get("/start")
async def start_wait_for_step_completion() -> str:
flow_id = required_query("workflowId")
await app_state.client.start_flow(
app_state.wait_for_step_completion,
flow_id,
JobSeekerData(1),
start_options(),
)
await app_state.client.wait_for_step_completion(
flow_id,
PERSIST_DATA_STEP,
WaitForStepCompletionOptions(),
)
persisted = await app_state.client.invoke_rpc(
app_state.wait_for_step_completion.get_job_seeker_data,
flow_id,
)
payload = json.dumps(asdict(persisted), sort_keys=True)
return f"success for workflow {flow_id} with data {payload}"

例子: examples/python/dex_examples/patterns/wait-for-step-completion/controller.py

相关内容