跳到主要内容

SubFlow

SubFlow 是 parent Flow 下的另一个 Flow execution。

SubFlow 有两个主要意义:

  1. 把复杂逻辑抽象到更小的 Flow 中。这样业务逻辑更容易维护,层级结构也更清晰。
  2. 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 FlowATTACHRESTART_IF_PREVIOUS_EXITS_ABNORMALLYALWAYS_RESTART
没有已有 FlowStartStartStart
RunningAttachAttachTerminate 后 start
Completed返回它的 result返回它的 resultStart
Failed、canceled、timed out 或 terminated返回它的 resultStartStart

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