SubFlow
SubFlow 是 parent Flow 下的另一个 Flow execution。
SubFlow 有两个主要意义:
- 把复杂逻辑抽象到更小的 Flow 中。这样业务逻辑更容易维护,层级结构也更清晰。
- Fan out 工作。一个 Flow execution 无法同时并发执行太多操作,但可以 delegate 给多个 SubFlow。每个 SubFlow 可以同时运行,所以 parent 可以达到更高的并发。
启动并等待 SubFlow
target Flow 必须注册在 Worker 上,而且它的 start Step 必须能接受传给 SubFlow 的 input。WaitFor 启动 child;child 完成后,Execute 读取结果。
class SubFlowParentStep(Step[int]):
def __init__(self, target: Flow[int]) -> None:
self.target = target
self.options = SubFlowOptions(
timeout=timedelta(hours=1),
timeout_policy=FlowTimeoutPolicy.CANCEL,
)
def wait_for(self, context: Context, input: int) -> Wait:
return Wait.until(SubFlow.run(self.target, input, self.options))
def execute(self, context: Context, input: int) -> StepDecision:
result = SubFlow.get_condition_results(context)
output = result.single_output(int)
return graceful_complete(f"{SubFlow.get_flow_id(context)}|{output}")
例子: examples/python/dex_examples/primitives/subflow/parent_flow.py
与其他 Condition 一起使用 SubFlow
使用 allOf,让 parent 等待所有 SubFlow 完成。
使用 anyOf 时,parent 不会等待所有 SubFlow。WaitFor 已经结束等待后,SubFlow 仍会继续运行。Execute 仍然可以拿到 SubFlow ID 来停止它,或者把 SubFlow ID 传给其他 Step,由它决定如何处理。
SubFlowOptions 和 SubFlowReusePolicy
SubFlowOptions 等同于 StartFlow API 上的 StartFlowOptions。大部分 option 都一样,只有 SubFlowReusePolicy 不同。
Dex Server 会根据 parent Flow ID、Step execution ID 和 condition index 生成 SubFlow ID。当 parent 通过 IDReusePolicy 复用 Flow ID,或者 time travel 时,它可能会进入同一个 Step execution,再尝试启动相同的 SubFlow ID。SubFlowReusePolicy 决定这种情况如何处理。
默认的 RESTART_IF_PREVIOUS_EXITS_ABNORMALLY 通常就是 parent 想要的:继续等待 running child,使用成功 child 的 result,重新启动异常结束的 child。
| 已存在的 child Flow | ATTACH | RESTART_IF_PREVIOUS_EXITS_ABNORMALLY | ALWAYS_RESTART |
|---|---|---|---|
| 没有已有 Flow | Start | Start | Start |
| Running | Attach | Attach | Terminate 后 start |
| Completed | 返回它的 result | 返回它的 result | Start |
| Failed、canceled、timed out 或 terminated | 返回它的 result | Start | Start |
Dex Server 会先比较为这次 child start 生成的 Request ID。Request ID 相同,就一定会 attach 到 running child 或返回它的 terminal result,和 ReusePolicy 无关。这样,Dex Server 接受了原始 start 但 parent 没收到 response 时,retry 也是安全的。Request ID 不同时,Dex Server 才根据表格应用 ReusePolicy。