高级:长生命周期父 Flow
保持有界的父 Step pool,同时让 Flow 持续接收新请求。
InitStep 将初始请求发布到 RequestChannel,把 Stopped 设为 false,并启动 N 个 HandleRequestStep。每个 worker 取得一条 FIFO 消息后跳转到 HandleSubFlowStep。后者启动并等待 ExampleSubFlow;完成后检查 Stopped,决定回到 HandleRequestStep 还是 graceful complete。
SendRequest 在发布前检查 Channel size。达到 buffer threshold 时返回 false,让调用方实施背压。Stop RPC 只负责把 Stopped 设为 true。
核心实现
class LongLiveHandleSubFlowStep(Step[str]):
def __init__(
self,
example_subflow: ExampleSubFlow,
stopped: Attribute[bool],
) -> None:
self.example_subflow = example_subflow
self.stopped = stopped
def get_step_type(self) -> str:
return "HandleSubFlowStep"
def wait_for(self, context: Context, request: str) -> Wait:
return Wait.until(SubFlow.run(self.example_subflow, request))
def execute(self, context: Context, request: str) -> StepDecision:
if self.stopped.get(context):
return graceful_complete()
return go_to(LongLiveHandleRequestStep, None)
例子: examples/python/dex_examples/patterns/parallel-subflows/advanced_long_live_parent_flow.py