动态并行 Step
InitStep 为每个输入项创建一个 DoWorkStep movement。Dex 并发调度这些 execution。每个 worker 会随机等待一小段时间,让并发完成顺序更直观。这个模式适用于并行度只能在运行时确定的场景。
核心实现
起始 Step 在运行时创建 movement 列表。worker 的随机延迟用来模拟耗时不同的工作。
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)
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/dynamic_parallel_steps_flow.py