高级:短生命周期父 Flow
排空当前工作后完成父 Flow,同时避免与新请求产生竞态。
HandleRequestStep.Execute 持有 CurrSubFlowNum 的 Attribute lock 并加一;HandleSubFlowStep.Execute 持有同一把 lock 并减一。
最后一个 active worker 使用 ForceCompleteIfChannelsEmpty。Dex 原子检查 RequestChannel:为空才完成;如果 RPC 已先发布新请求,则再次调度 HandleRequestStep。这里不能用应用代码先查 size 再完成,否则两次操作之间会产生竞态。
核心实现
class ShortLiveHandleSubFlowStep(Step[str]):
def __init__(
self,
example_subflow: ExampleSubFlow,
request_channel: Channel[str],
curr_subflow_num: Attribute[int],
) -> None:
self.example_subflow = example_subflow
self.request_channel = request_channel
self.curr_subflow_num = curr_subflow_num
def get_step_type(self) -> str:
return "HandleSubFlowStep"
def get_step_options(self) -> StepOptions:
return StepOptions(execute_lock_attributes=(self.curr_subflow_num.lock(),))
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:
current = (self.curr_subflow_num.get(context) or 0) - 1
self.curr_subflow_num.set(context, current)
if current == 0:
return force_complete_if_channels_empty(
None,
StepMovement.of(ShortLiveHandleRequestStep, None),
self.request_channel,
)
return go_to(ShortLiveHandleRequestStep, None)
例子: examples/python/dex_examples/patterns/parallel-subflows/advanced_short_live_parent_flow.py