跳转到内容

工作流的并发执行

除了循环、分支和流式处理,工作流还可以并发运行步骤。当您有多个可以相互独立运行的步骤,并且它们包含耗时的操作时,这非常有用,它们 await,允许其他步骤并行运行。

要触发多个步骤而发出多个事件,您可以使用 ctx.send_event()

import asyncio
from workflows import Workflow, Context, step
from 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 个 StepTwoEventstep_two 步骤使用了 num_workers=4 装饰器,这告诉工作流最多可以同时运行 4 个该步骤的实例(这是默认设置)。

如果你执行前面的示例,你会注意到工作流会在任意一个查询首先完成后停止。有时这很有用,但其他时候你可能希望等待所有耗时操作完成后再继续下一步。你可以使用 collect_events 来实现:

import asyncio
from workflows import Workflow, Context, step
from 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 asyncio
from workflows import Workflow, Context, step
from 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 的顺序返回。

这个工作流的可视化效果相当令人满意:

A concurrent workflow