首个获胜的并行 Step
InitStep 启动多个 DoWorkStep execution。每个 worker 随机等待一小段时间,从而形成真正的竞速。最先返回 decision 的 execution 会取消同类的 sibling execution。Flow 保留第一个结果,并停止已经无法影响结果的工作。
核心实现
取消选择器只匹配同一 Step 类型的 sibling execution,不会影响无关分支。
class DoWorkStep(Step[int]):
async def execute( # type: ignore[override]
self, context: AsyncContext, input: int
) -> StepDecision:
await asyncio.sleep(random.uniform(0.05, 0.5))
return graceful_complete(input).with_canceling_sibling_steps(DoWorkStep)
class InitStep(Step[int]):
def execute(self, context: Context, input: int) -> StepDecision:
return go_to_many(
*(StepMovement.of(DoWorkStep, index) for index in range(input))
)
例子: examples/python/dex_examples/patterns/parallel/first_win_parallel_steps_flow.py