工作流的并发执行
除了循环、分支和流式处理,工作流还可以并发运行步骤。当您有多个可以相互独立运行的步骤,并且它们包含耗时的操作时,这非常有用,它们 await,允许其他步骤并行运行。
要触发多个步骤而发出多个事件,您可以使用 ctx.send_event():
import asynciofrom workflows import Workflow, Context, stepfrom workflows.events import Event, StartEvent, StopEvent
class StepTwoEvent(Event): query: str
class ParallelFlow(Workflow): @step async def start(self, ctx: Context, ev: StartEvent) -> StepTwoEvent | None: ctx.send_event(StepTwoEvent(query="Query 1")) ctx.send_event(StepTwoEvent(query="Query 2")) ctx.send_event(StepTwoEvent(query="Query 3"))
@step(num_workers=4) async def step_two(self, ev: StepTwoEvent) -> StopEvent: print("Running slow query ", ev.query) await asyncio.sleep(random.randint(0, 5))
return StopEvent(result=ev.query)在这个示例中,我们的 start 步骤生成了 3 个 StepTwoEvent。 step_two 步骤使用了 num_workers=4 装饰器,这告诉工作流最多可以同时运行 4 个该步骤的实例(这是默认设置)。
如果你执行前面的示例,你会注意到工作流会在任意一个查询首先完成后停止。有时这很有用,但其他时候你可能希望等待所有耗时操作完成后再继续下一步。你可以使用 collect_events 来实现:
import asynciofrom workflows import Workflow, Context, stepfrom workflows.events import Event, StartEvent, StopEvent
class StepTwoEvent(Event): query: str
class StepThreeEvent(Event): result: str
class ConcurrentFlow(Workflow): @step async def start(self, ctx: Context, ev: StartEvent) -> StepTwoEvent | None: ctx.send_event(StepTwoEvent(query="Query 1")) ctx.send_event(StepTwoEvent(query="Query 2")) ctx.send_event(StepTwoEvent(query="Query 3"))
@step(num_workers=4) async def step_two(self, ctx: Context, ev: StepTwoEvent) -> StepThreeEvent: print("Running query ", ev.query) await asyncio.sleep(random.randint(1, 5)) return StepThreeEvent(result=ev.query)
@step async def step_three( self, ctx: Context, ev: StepThreeEvent ) -> StopEvent | None: # wait until we receive 3 events result = ctx.collect_events(ev, [StepThreeEvent] * 3) if result is None: return None
# do something with all 3 results together print(result) return StopEvent(result="Done")collect_events 方法位于 Context 上,它接收触发步骤的事件以及要等待的事件类型数组。在本例中,我们正在等待3个相同 StepThreeEvent 类型的事件。
每当接收到一个 StepThreeEvent 时,step_three 步骤都会被触发,但 collect_events 将返回 None,直到所有3个事件都已接收。此时,该步骤将继续执行,您可以同时对这3个结果进行操作。
从 collect_events 返回的 result 是一个按接收顺序排列的已收集事件数组。
当然,您无需等待同类型事件。您可以等待任意组合的事件,如以下示例所示:
import asynciofrom workflows import Workflow, Context, stepfrom workflows.events import Event, StartEvent, StopEvent
class StepAEvent(Event): query: str
class StepBEvent(Event): query: str
class StepCEvent(Event): query: str
class StepACompleteEvent(Event): result: str
class StepBCompleteEvent(Event): result: str
class StepCCompleteEvent(Event): result: str
class ConcurrentFlow(Workflow): @step async def start( self, ctx: Context, ev: StartEvent ) -> StepAEvent | StepBEvent | StepCEvent | None: ctx.send_event(StepAEvent(query="Query 1")) ctx.send_event(StepBEvent(query="Query 2")) ctx.send_event(StepCEvent(query="Query 3"))
@step async def step_a(self, ctx: Context, ev: StepAEvent) -> StepACompleteEvent: print("Doing something A-ish") return StepACompleteEvent(result=ev.query)
@step async def step_b(self, ctx: Context, ev: StepBEvent) -> StepBCompleteEvent: print("Doing something B-ish") return StepBCompleteEvent(result=ev.query)
@step async def step_c(self, ctx: Context, ev: StepCEvent) -> StepCCompleteEvent: print("Doing something C-ish") return StepCCompleteEvent(result=ev.query)
@step async def step_three( self, ctx: Context, ev: StepACompleteEvent | StepBCompleteEvent | StepCCompleteEvent, ) -> StopEvent: print("Received event ", ev.result)
# wait until we receive 3 events if ( ctx.collect_events( ev, [StepCCompleteEvent, StepACompleteEvent, StepBCompleteEvent], ) is None ): return None
# do something with all 3 results together return StopEvent(result="Done")我们进行了几项更改以处理多种事件类型:
startstart 现在被声明为可发出3种不同的事件类型step_threestep_three 现在被声明为接受3种不同的事件类型collect_eventscollect_events 现在接受一个事件类型数组作为等待条件
请注意,传递给 collect_events 的事件类型数组中的顺序很重要。无论事件何时被接收,它们都将按照传递给 collect_events 的顺序返回。
这个工作流的可视化效果相当令人满意:
