跳到主要内容

等待并行 Step

Definition graph

AwaitParallelStepsFlow

Valid

python · examples/python/dex_examples/patterns/parallel/await_parallel_steps_flow.py

InitStep 同时启动 DoWorkStepAwaitStep。每个 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

相关内容