等待并行 Step
InitStep 同时启动 DoWorkStep 和 AwaitStep。每个 worker 随机等待一小段时间,然后向 CompleteCh 发布一条消息。AwaitStep 等到所需数量的消息后才完成 Flow。Channel 消息是持久化的,因此即使消息先于 AwaitStep 开始等待,也仍然会被计入。
核心实现
协调 Step 与 worker 一起启动,并等待每个 worker 各发布一条完成消息。
class DoWorkStep(Step[int]):
def __init__(self, complete_ch: Channel[None]) -> None:
self.complete_ch = complete_ch
async def execute( # type: ignore[override]
self, context: AsyncContext, input: int
) -> StepDecision:
await asyncio.sleep(random.uniform(0.05, 0.5))
self.complete_ch.publish(context, None)
return dead_end()
class AwaitStep(Step[int]):
def __init__(self, complete_ch: Channel[None]) -> None:
self.complete_ch = complete_ch
def wait_for(self, context: Context, input: int) -> Wait:
return Wait.until(self.complete_ch.for_n(input))
def execute(self, context: Context, input: int) -> StepDecision:
return graceful_complete(input)
class InitStep(Step[int]):
def execute(self, context: Context, input: int) -> StepDecision:
movements: list[StepMovement[Any]] = [StepMovement.of(AwaitStep, input)]
movements.extend(
StepMovement.of(DoWorkStep, index) for index in range(input)
)
return go_to_many(*movements)
例子: examples/python/dex_examples/patterns/parallel/await_parallel_steps_flow.py