워크플로 동시 실행(Concurrent Execution)
워크플로 동시 실행(Concurrent Execution)
워크플로는 여러 스텝을 동시에 실행할 수 있어요. 여러 스텝이 서로 독립적이고 각자 느린 작업을 기다리는 상황이라면 병렬로 실행하는 게 더 빠릅니다.
가장 흔한 패턴은 팬아웃(fan-out)과 팬인(fan-in)이에요. 작업을 여러 조각으로 나누어 동시에 실행한 다음, 결과를 다시 합치죠. 이걸 스텝 시그니처에 직접 작성합니다. 스텝에서 list를 반환하면 요소마다 하나의 이벤트로 팬아웃되고, list 파라미터를 받으면 전체 배치를 한 번에 모아 팬인돼요. @step 데코레이터가 이 타입들을 읽습니다. 검증기와 시각화 도구는 여러분이 추가로 뭘 하지 않아도 각 생성 스텝과 그 이벤트를 소비하는 스텝을 자동으로 연결해 줘요. 시그니처에 따르지 않는 이벤트를 내보내야 할 때는 ctx.send_event로 직접 보낼 수 있습니다. 이 페이지 끝의 동적 API 섹션이 그 내용을 다룹니다.
동시 워크플로 실행 수 제한하기
워크플로 전체의 동시성과 스텝 단위의 동시성은 별개예요. num_concurrent_runs를 설정하면 한 번에 활성화될 수 있는 run() 호출 수를 제한합니다.
workflow = ParallelFlow(num_concurrent_runs=4)
값은 양의 정수이거나 None이어야 해요. 기본값은 None으로 무제한 실행을 허용합니다. 이 제한은 프로세스 내부의 동시성을 제한하며, 추가 호출은 활성 실행이 끝날 때까지 기다려요. DBOS 기반 워크플로도 워커별로 적용되는 같은 기능을 지원합니다.
각 실행 내부의 동시성을 다루려면 @step(num_workers=...)을 써요. 이 값은 그 스텝이 얼마나 많은 복사본으로 이벤트를 동시에 처리할 수 있는지를 정하며, 워크플로 실행 횟수를 제한하지는 않습니다.
실행들에 걸친 활성 스텝 제한하기
모든 실행에 걸쳐 한 스텝의 동시 복사본 수를 제한하고 싶다면(예: 외부 API에 대한 동시 호출 상한), 스텝 사이에서 asyncio.Semaphore를 공유하면 됩니다. 스텝은 평범한 async 함수이므로 라이브러리 지원이 필요 없어요.
import asyncio
API_SLOTS = asyncio.Semaphore(3)
class Fetcher(Workflow):
@step
async def fetch(self, ev: FetchRequest) -> FetchResult:
async with API_SLOTS:
data = await call_external_api(ev.url)
return FetchResult(data=data)
모든 실행이 하나의 세마포어를 공유하므로, 어떤 실행에 속하든 어느 순간에도 fetch 스텝은 기껏해야 3개만 API 호출 중이에요. 세마포어는 하나의 Python 프로세스 안에 살아있고, 여러 프로세스로 배포하면 각각 자기 것을 갖습니다.
팬아웃: 리스트 반환하기
스텝에서 list를 반환하면 각 요소가 자기만의 이벤트로 발화됩니다. 아래 예에서는 work 아래에서 다섯 개의 Task가 동시에 실행돼요.
import asyncio
import random
from workflows import Workflow, step
from workflows.events import Event, StartEvent, StopEvent
class Task(Event):
n: int
class Done(Event):
n: int
class ParallelFlow(Workflow):
@step
async def fan_out(self, ev: StartEvent) -> list[Task]:
return [Task(n=i) for i in range(5)]
@step(num_workers=5)
async def work(self, ev: Task) -> Done:
await asyncio.sleep(random.randint(0, 5))
return Done(n=ev.n)
전체 리스트가 하나의 배치예요. num_workers=5는 work 스텝의 복사본 최대 5개를 동시에 실행하게 합니다. []를 반환하면 아무것도 발화되지 않지만 스텝은 여전히 완료되고 배치는 즉시 닫힙니다.
팬인: 리스트 받기
단일 이벤트 대신 list 파라미터를 받으면 스텝이 전체 배치를 모은 뒤, 모든 것이 담긴 채 한 번만 발화돼요.
class ConcurrentFlow(Workflow):
@step
async def fan_out(self, ev: StartEvent) -> list[Task]:
return [Task(n=i) for i in range(5)]
@step(num_workers=5)
async def work(self, ev: Task) -> Done:
await asyncio.sleep(random.randint(1, 5))
return Done(n=ev.n)
@step
async def join(self, events: list[Done]) -> StopEvent:
return StopEvent(result=sorted(e.n for e in events))
fan_out은 list[Task]를 반환하고 join은 list[Done]을 받으므로, 프레임워크는 타입만으로 배치를 알고 join이 완료되면 발화시킵니다.
리스트의 이벤트는 완료 순서, 즉 각 워커가 끝난 순서예요. 이건 fan_out이 보낸 순서와 다릅니다. 고정된 순서가 필요하다면 join이 여기서 sorted로 하는 것처럼 스스로 정렬하세요.
워커는 None을 반환해 자기 분기를 버릴 수 있어요. 프레임워크는 여전히 값을 반환하는 분기를 추적하므로 join은 그 분기들만 담아 한 번 발화됩니다.
@step(num_workers=5)
async def work(self, ev: Task) -> Done | None:
if ev.n % 2 == 0:
return None # drop this branch
return Done(n=ev.n)
여기서 짝수 번호의 작업은 None을 반환하므로 join은 홀수 Done 이벤트만 받아요.
일찍 릴리즈하기
기본적으로 list 조인은 전체 배치를 기다립니다. 첫 결과에만 반응하고 싶다면 파라미터를 Collect로 감싸세요. Take(n)은 n번째 도착 시 첫 n개 이벤트로 발화됩니다. 처음 몇 개 결과만, 또는 가장 먼저 끝난 하나만 필요할 때 쓰는 방식이에요.
from typing import Annotated
from workflows.collect import Collect, Take
class FastestWins(Workflow):
@step
async def fan_out(self, ev: StartEvent) -> list[Task]:
return [Task(n=i) for i in range(5)]
@step(num_workers=5)
async def work(self, ev: Task) -> Done:
await asyncio.sleep(random.randint(1, 5))
return Done(n=ev.n)
@step
async def first(
self, events: Annotated[list[Done], Collect(Take(1))]
) -> StopEvent:
return StopEvent(result=events[0].n)
Take(1)은 어떤 작업이든 가장 먼저 끝나는 것에서 멈춥니다. 다른 작업들은 계속 실행돼요. 아무것도 취소하지 않고, 조인에 도달하지 못할 뿐입니다. 평범한 list[Done] 파라미터는 Annotated[list[Done], Collect(All())]과 같아요. 인자 없는 Collect()는 코드에서 검색하기 쉽게 같은 기본값을 씁니다.
혼합 이벤트 타입 팬인
배치가 반드시 하나의 이벤트 타입일 필요는 없어요. list[A | B]는 두 타입을 평평하게 섞은 배치를 모으고, 스텝은 그중 하나를 받습니다.
@step
async def join(
self, events: list[StepACompleteEvent | StepBCompleteEvent]
) -> StopEvent:
...
리스트 대신 각 타입을 하나씩 기다리려면 스텝에 이벤트마다 파라미터 하나씩을 주세요. 각 파라미터가 이벤트를 받으면 스텝이 한 번 발화됩니다. 각 파라미터는 타입으로 매칭돼요.
from workflows import Workflow, step
from workflows.events import Event, StartEvent, StopEvent
class StepACompleteEvent(Event):
result: str
class StepBCompleteEvent(Event):
result: str
class StepCCompleteEvent(Event):
result: str
class ConcurrentFlow(Workflow):
@step
async def step_a(self, ev: StartEvent) -> StepACompleteEvent:
return StepACompleteEvent(result="Query 1")
@step
async def step_b(self, ev: StartEvent) -> StepBCompleteEvent:
return StepBCompleteEvent(result="Query 2")
@step
async def step_c(self, ev: StartEvent) -> StepCCompleteEvent:
return StepCCompleteEvent(result="Query 3")
@step
async def assemble(
self,
a: StepACompleteEvent,
b: StepBCompleteEvent,
c: StepCCompleteEvent,
) -> StopEvent:
return StopEvent(result=[a.result, b.result, c.result])
step_a, step_b, step_c는 모두 같은 StartEvent에서 실행되므로 동시에 실행돼요. assemble은 각각에 대해 파라미터를 하나씩 갖고 있으므로 세 개가 모두 도착하면 한 번 발화됩니다.
중첩(Nesting)
팬아웃을 중첩할 수 있어요. 팬아웃 안에서 또 팬아웃하면 중첩 배치가 됩니다. 내부 조인은 바깥 요소마다 한 번씩 발화되고, 그다음 외부 조인이 모든 내부 결과를 모아 한 번 발화돼요.
class InnerTask(Event):
outer: int
inner: int
class InnerDone(Event):
outer: int
inner: int
class InnerSummary(Event):
outer: int
total: int
class Nested(Workflow):
@step
async def outer(self, ev: StartEvent) -> list[Task]:
return [Task(n=o) for o in range(3)]
@step
async def inner(self, ev: Task) -> list[InnerTask]:
return [InnerTask(outer=ev.n, inner=i) for i in range(2)]
@step
async def inner_work(self, ev: InnerTask) -> InnerDone:
return InnerDone(outer=ev.outer, inner=ev.inner)
@step
async def per_inner(self, events: list[InnerDone]) -> InnerSummary:
return InnerSummary(outer=events[0].outer, total=len(events))
@step
async def per_outer(self, events: list[InnerSummary]) -> StopEvent:
return StopEvent(result=sorted((s.outer, s.total) for s in events))
각 조인은 자기 수준에 머물러요. per_inner는 외부 Task마다 한 번씩 총 3번 실행되고, per_outer는 세 요약을 모아 한 번 실행됩니다.
동적 API
스텝이 반환하기 전에 이벤트를 내보내야 할 때는 ctx.send_event로 직접 보내고 ctx.collect_events로 모아요. 조건에 따라 일부 이벤트만 방출하거나, 방출할 이벤트 수를 미리 모를 때, 또는 프로듀서가 여전히 실행 중인 동안 다운스트림 작업을 시작하고 싶을 때 씁니다.
여기엔 트레이드오프가 있어요. list 반환은 올인(전부 또는 없음) 방식이에요. 스텝이 반환할 때 하나의 배치로 나가므로, 스텝이 먼저 오류를 일으키면 아무것도 방출되지 않습니다. 반면 ctx.send_event는 호출하는 즉시 발화돼요. 다운스트림 스텝이 프로듀서 스텝이 끝나기 전에 시작할 수 있지만, 이미 보낸 것은 나중에 스텝이 실패해도 밖에 남아 있어요.
ctx.send_event는 한 번에 하나의 이벤트를 방출합니다.
import asyncio
import random
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는 이벤트를 반환하는 대신 ctx.send_event로 방출합니다. 함수가 None을 반환해도 반환 어노테이션에 StepTwoEvent가 여전히 포함되어 있어서, 검증과 다이어그램이 이 스텝이 그 이벤트를 만들 수 있음을 알아요. 보낸 이벤트를 시그니처에서 빼도 런타임은 여전히 보낼 수 있지만, 정적 검증과 시각화는 그 간선을 추론하지 못합니다.
여러 개의 수동 전송 이벤트를 기다리려면 ctx.collect_events를 써요.
import asyncio
import random
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:
# 3개의 이벤트를 받을 때까지 대기
result = ctx.collect_events(ev, [StepThreeEvent] * 3)
if result is None:
return None
# 3개 결과를 함께 처리
print(result)
return StopEvent(result="Done")
ctx.collect_events는 트리거 이벤트와 기다릴 타입 리스트를 받아요. step_three는 모든 StepThreeEvent에서 실행되지만, collect_events는 세 개가 모두 도착할 때까지 None을 반환합니다. 그다음 도착 순서대로 리스트로 돌려줘요. 기대할 이벤트 수를 직접 추적해야 하며, 그게 여기서의 3이에요.
어떤 타입 조합이든 기다릴 수 있으며 반드시 한 타입만 반복할 필요는 없어요. 넘겨준 순서대로, 실제로 언제 도착했는지와 무관하게 같은 순서로 돌아옵니다.
@step
async def step_three(
self,
ctx: Context,
ev: StepACompleteEvent | StepBCompleteEvent | StepCCompleteEvent,
) -> StopEvent | None:
if (
ctx.collect_events(
ev,
[StepCCompleteEvent, StepACompleteEvent, StepBCompleteEvent],
)
is None
):
return None
return StopEvent(result="Done")
팬아웃을 내구성 있게 만들기
긴 팬아웃은 체크포인팅에 잘 맞아요. 대기 중인 이벤트와 부분 팬인 상태를 직렬화해서 재시작 후 다시 이어갈 수 있죠. 체크포인트 루프와 예제는 Writing durable workflows를 참고하세요.
더 알아보기
- 분기와 루프(Branches and loops) — 조건 분기와 반복 로직 기초.
- 스테이트 관리(Managing State) — 동시 실행 중 공유 상태 다루기.