에러 처리(Error Handling)와 재시도
에러 처리(Error Handling)와 재시도
스텝이 실행에 실패하면 워크플로 전체가 실패할 수 있어요. 하지만 오류가 예상되는 경우도 많고, 안전하게 재시도할 수 있는 경우도 많습니다. 네트워크의 일시적인 혼잡으로 타임아웃된 HTTP 요청, 또는 레이트 리미터에 걸린 외부 API 호출을 떠올려 보세요.
스텝을 다시 시도하고 싶은 모든 상황에서 **재시도 정책(retry policy)**을 쓸 수 있어요. 재시도 정책은 워크플로가 스텝을 여러 번 실행하도록 지시하고, 각 시도 사이에 기다릴 시간, 어떤 오류를 재시도할지, 언제 포기할지를 제어합니다.
재시도는 같은 입력 이벤트로 스텝 본문 전체를 다시 실행해요. 실패하는 호출 이전의 부수 효과(side effect)는 멱등(idempotent)하게 유지하거나, 재시도 가능한 작업 뒤로 옮기세요. 스텝이 예외를 던지기 전에 상태를 쓰거나 외부 서비스를 호출했다면 그 작업은 다시 일어날 수 있어요.
재시도 모듈은 조합 가능한 세 종류의 빌딩 블록으로 구성됩니다.
- 재시도 조건(Retry conditions) — 예외가 재시도 가능한지 결정한다.
- 대기 전략(Wait strategies) — 다음 시도 전에 얼마나 잠들지 결정한다.
- 중단 조건(Stop conditions) — 워크플로가 언제 포기할지 결정한다.
retry_policy(retry=..., wait=..., stop=...)로 정책에 조합하고 결과를 @step 데코레이터에 넘기세요. 재시도 조건은 |와 &를, 대기 전략은 +를, 중단 조건은 |와 &를 지원합니다.
기본 사용법
from workflows import Workflow, Context, step
from workflows.events import StartEvent, StopEvent
from workflows.retry_policy import (
retry_policy,
stop_after_attempt,
wait_fixed,
)
class MyWorkflow(Workflow):
@step(
retry_policy=retry_policy(
wait=wait_fixed(5),
stop=stop_after_attempt(10),
)
)
async def flaky_step(self, ctx: Context, ev: StartEvent) -> StopEvent:
result = flaky_call() # this might raise
return StopEvent(result=result)
인자 없이 retry_policy()는 모든 예외를 최대 3번 시도하고 각 사이에 5초 고정 대기를 합니다. retry=, wait=, stop=로 어느 구성 요소든 바꿀 수 있어요.
예외 필터링하기
재시도 조건으로 어떤 예외를 재시도할지 제어하세요.
from workflows import Workflow, step
from workflows.events import StartEvent, StopEvent
from workflows.retry_policy import (
retry_policy,
retry_if_exception_message,
retry_if_exception_type,
stop_after_attempt,
stop_before_delay,
wait_fixed,
wait_random,
)
class MyWorkflow(Workflow):
@step(
retry_policy=retry_policy(
retry=retry_if_exception_type((TimeoutError, ConnectionError))
| retry_if_exception_message(match="rate limit|temporarily unavailable"),
wait=wait_fixed(1) + wait_random(0, 1),
stop=stop_after_attempt(5) | stop_before_delay(30),
)
)
async def call_provider(self, ev: StartEvent) -> StopEvent:
result = await flaky_call()
return StopEvent(result=result)
보다 명시적인 스타일을 선호한다면 이름 붙은 조합기가 동일하게 동작합니다.
from workflows.retry_policy import (
retry_policy,
retry_any,
retry_if_exception_message,
retry_if_exception_type,
stop_any,
stop_after_attempt,
stop_before_delay,
wait_combine,
wait_fixed,
wait_random,
)
policy = retry_policy(
retry=retry_any(
retry_if_exception_type((TimeoutError, ConnectionError)),
retry_if_exception_message(match="rate limit|temporarily unavailable"),
),
wait=wait_combine(wait_fixed(1), wait_random(0, 1)),
stop=stop_any(stop_after_attempt(5), stop_before_delay(30)),
)
지수 백오프(Exponential backoff)
LLM 프로바이더나 레이트 제한 API를 호출하는 스텝에는 지수 백오프가 썬더링 허드(thundering-herd) 효과를 피하게 해 줍니다.
from workflows import Workflow, Context, step
from workflows.events import StartEvent, StopEvent
from workflows.retry_policy import (
retry_policy,
stop_after_attempt,
wait_exponential_jitter,
)
class MyWorkflow(Workflow):
@step(
retry_policy=retry_policy(
wait=wait_exponential_jitter(initial=1, exp_base=2, max=30, jitter=1),
stop=stop_after_attempt(5),
)
)
async def call_llm(self, ctx: Context, ev: StartEvent) -> StopEvent:
result = await llm_call() # this might raise on rate-limit
return StopEvent(result=result)
커스텀 재시도 정책
조합 가능한 API가 유스케이스를 다루지 못한다면 커스텀 정책을 작성할 수 있어요. 유일한 요구사항은 RetryPolicy 프로토콜과 일치하는 next 메서드를 가진 클래스입니다.
def next(
self,
elapsed_time: float,
attempts: int,
error: Exception,
*,
seed: int | None = None,
) -> float | None:
...
재시도 전에 기다릴 초 수를 반환하거나, 멈추려면 None을 반환하세요. attempts는 지금까지의 실패 횟수로 첫 실패에서 1부터 시작합니다. 선택적 seed는 내구성 있는 런타임이 리플레이 중에 지터를 결정적으로 만들기 위해 사용해요.
예를 들어 이 정책은 금요일에만 재시도합니다.
from datetime import datetime
class RetryOnFridayPolicy:
def next(
self,
elapsed_time: float,
attempts: int,
error: Exception,
*,
seed: int | None = None,
) -> float | None:
if datetime.today().strftime("%A") == "Friday":
return 5 # retry in 5 seconds
return None # don't retry
지원 중단된 편의 생성자(Deprecated convenience constructors)
주의: ConstantDelayRetryPolicy와 ExponentialBackoffRetryPolicy는 조합 가능한 API보다 앞서 만들어진 것으로, 하위 호환을 위해서만 유지됩니다. 명시적인 retry/wait/stop 인자를 쓰는 retry_policy(...)를 선호하세요.
from workflows.retry_policy import ConstantDelayRetryPolicy
# Deprecated, equivalent to:
# retry_policy(wait=wait_fixed(5), stop=stop_after_attempt(10))
policy = ConstantDelayRetryPolicy(delay=5, maximum_attempts=10)
from workflows.retry_policy import ExponentialBackoffRetryPolicy
# Deprecated, equivalent to:
# retry_policy(wait=wait_random_exponential(multiplier=1, exp_base=2, max=30),
# stop=stop_after_attempt(5))
policy = ExponentialBackoffRetryPolicy(
initial_delay=1, multiplier=2, max_delay=30, maximum_attempts=5,
)
스텝 안에서 재시도 상태 살펴보기
Context.retry_info()는 현재 시도를 설명하는 RetryInfo를 돌려줘요. 재시도 사이에 동작을 바꾸는 데 씁니다. 검색을 넓히거나, 타임아웃을 줄이거나, 경고를 로깅하는 식이죠.
from workflows import Context, Workflow, step
from workflows.events import StartEvent, StopEvent
from workflows.retry_policy import retry_policy, stop_after_attempt
class MyWorkflow(Workflow):
@step(retry_policy=retry_policy(stop=stop_after_attempt(3)))
async def flaky(self, ctx: Context, ev: StartEvent) -> StopEvent:
info = ctx.retry_info()
if info.last_exception is not None:
logger.warning(
"retry %d after %s", info.retry_number, str(info.last_exception)
)
...
retry_number는 첫 실행에서 0이에요. last_exception과 last_failed_at은 첫 실패 전까지 None이며, 이후 가장 최근의 이전 실패를 설명합니다.
재시도가 소진됐을 때 복구하기
정책이 포기하면 예외가 전파되어 워크플로가 실패합니다. 단, @catch_error 핸들러가 선언되어 있다면 실패하지 않아요. 핸들러는 StepFailedEvent를 받는 스텝입니다. StopEvent를 반환해 실행을 우아하게 끝내거나, 다른 이벤트를 반환해 워크플로를 재라우팅하거나, 예외를 던져 새 오류로 실패시킬 수 있어요. ev.step_name과 ev.exception을 검사해 결정하세요.
from workflows import Context, Workflow, catch_error, step
from workflows.events import StartEvent, StepFailedEvent, StopEvent
from workflows.retry_policy import retry_policy, stop_after_attempt, wait_fixed
class Pipeline(Workflow):
@step(retry_policy=retry_policy(wait=wait_fixed(1), stop=stop_after_attempt(3)))
async def fetch(self, ev: StartEvent) -> StopEvent:
return StopEvent(result=await call_api())
@catch_error
async def on_failure(
self, ctx: Context, ev: StepFailedEvent
) -> StopEvent:
return StopEvent(
result={"failed_step": ev.step_name, "error": str(ev.exception)}
)
맨몸의 @catch_error는 와일드카드예요. 재시도가 소진된 스텝 중 더 구체적인 핸들러에 속하지 않은 것을 모두 잡습니다. 와일드카드는 워크플로당 하나만 허용됩니다.
특정 스텝으로 범위 지정하기
for_steps=[...]로 핸들러가 다루는 스텝을 제한하세요. 스텝 이름은 범위 지정 핸들러 한 곳에만 나타날 수 있고, 알 수 없는 이름은 생성 시점에 거부됩니다.
class Pipeline(Workflow):
@step(retry_policy=...)
async def fetch(self, ev: StartEvent) -> FetchedEvent: ...
@step(retry_policy=...)
async def parse(self, ev: FetchedEvent) -> StopEvent: ...
@catch_error(for_steps=["fetch"])
async def on_fetch_failure(
self, ctx: Context, ev: StepFailedEvent
) -> StopEvent:
return StopEvent(result={"fallback": True})
@catch_error # wildcard; covers `parse` and anything else
async def on_any_failure(
self, ctx: Context, ev: StepFailedEvent
) -> StopEvent:
return StopEvent(result={"aborted": ev.step_name})
복구 예산(Recovery budget)
핸들러는 그래프로 흘러들어가는 이벤트를 스스로 방출할 수 있으므로, 한 혈통(lineage)이 같은 핸들러를 다시 들어갈 수 있어요. max_recoveries(기본 1)는 이벤트 혈통별로 그런 일이 몇 번까지 일어날 수 있는지를 제한하고, 그보다 많으면 재진입 대신 워크플로가 실패합니다. 유한 재시도 루프에 참여하는 핸들러라면 더 높게, 종료 핸들러라면 1로 유지하세요.
@catch_error(for_steps=["fetch", "retry_fetch"], max_recoveries=2)
async def on_failure(self, ctx: Context, ev: StepFailedEvent) -> RetryFetch:
return RetryFetch()
더 알아보기
- 이벤트 스트리밍(Streaming events) — 재시도 소진 후의 실패 이벤트(
WorkflowFailedEvent) 처리. - 비동기 워크플로(Async Workflows) — 실패를 다시 시도할 때 이벤트 루프를 막지 않는 법.