워크플로 상태 관리(Managing State)

워크플로 상태 관리(Managing State)

워크플로 실행 하나마다 Context가 있고, 각 컨텍스트에는 상태 저장소(state store)가 있어요. 실행 중에 스텝들이 공유해야 하는 값, 또는 같은 컨텍스트를 재사용·복원해서 다음 실행으로 의도적으로 이어가고 싶은 값을 보관하는 데 씁니다.

상태는 무거운 클라이언트, 인덱스, 파일 핸들, 그 외 런타임 의존성을 위한 자리가 아니에요. 그런 것들은 리소스(resources)에 두세요. 상태는 워크플로를 스냅샷하거나 재개할 때 직렬화해도 괜찮은 데이터여야 합니다.

기본적으로 워크플로는 타입이 없는(untyped) 상태 저장소를 초기화합니다. ctx.store를 통해 읽고 쓸 수 있어요.

from workflows import Workflow, Context, step
from workflows.events import StartEvent, StopEvent


class MyWorkflow(Workflow):

    @step
    async def my_step(self, ctx: Context, ev: StartEvent) -> StopEvent:
        current_count = await ctx.store.get("count", default=0)
        current_count += 1
        await ctx.store.set("count", current_count)
        return StopEvent(result=current_count)

출처: 공식문서 - Managing State

상태 잠그기(Locking the State)

여러 스텝이 동시에 상태를 조작하는 경우가 있어요. 그럴 때는 edit_state()로 상태를 잠가 업데이트가 원자적(atomic)이 되게 하세요.

@step
async def my_step(self, ctx: Context, ev: StartEvent) -> StopEvent:
    async with ctx.store.edit_state() as state:
        current_count = state.get("count", 0)
        state["count"] = current_count + 1
        result = state["count"]
    return StopEvent(result=result)

블록이 실행되는 동안 다른 스텝은 상태를 수정할 수 없어요. 블록은 작게 유지하세요. 읽고-수정하고-쓰는 작업을 블록 안에 넣고, 느린 LLM이나 네트워크 호출은 블록 밖에 두는 겁니다.

타입이 있는 상태 추가하기(Adding Typed State)

워크플로 상태의 형태가 정해져 있는 경우가 많아요. 그럴 땐 Pydantic 모델을 쓰세요. 그러면 다음과 같은 이점이 생깁니다.

  • 상태에 대한 타입 힌트를 얻는다
  • 상태의 자동 검증을 얻는다
  • (선택적으로) validatorsserializers로 상태의 직렬화·역직렬화를 완전히 제어한다

참고: 모든 필드에 기본값이 있는 Pydantic 모델을 쓰세요. 그래야 Context가 상태를 자동으로 초기화할 수 있어요.

워크플로 + Pydantic을 활용하는 간단한 예제입니다.

from pydantic import BaseModel, Field


class CounterState(BaseModel):
    count: int = Field(default=0)

그다음 워크플로 상태를 상태 모델로 어노테이션하면 됩니다.

from workflows import Workflow, Context, step
from workflows.events import (
    StartEvent,
    StopEvent,
)


class MyWorkflow(Workflow):
    @step
    async def start(
        self,
        ctx: Context[CounterState], ev: StartEvent
    ) -> StopEvent:
        # Allows for atomic state updates
        async with ctx.store.edit_state() as state:
            state.count += 1

        return StopEvent(result="Done!")

타입이 있는 상태를 필드 하나씩 다룰 수도 있어요.

@step
async def start(self, ctx: Context[CounterState], ev: StartEvent) -> StopEvent:
    current = await ctx.store.get("count")
    await ctx.store.set("count", current + 1)
    state = await ctx.store.get_state()
    return StopEvent(result=state.count)

실행들에 걸쳐 컨텍스트 유지하기(Maintaining Context Across Runs)

워크플로의 여러 실행에 걸쳐 상태를 유지하고 싶다면 컨텍스트를 만들고 같은 것을 .run()에 넘기세요.

workflow = MyWorkflow()
ctx = Context(workflow)

handler = workflow.run(ctx=ctx)
result = await handler

# Optional: save the ctx somewhere and restore
# ctx_dict = ctx.to_dict()
# ctx = Context.from_dict(workflow, ctx_dict)

# continue with next run
handler = workflow.run(ctx=ctx)
result = await handler

컨텍스트가 아직 실행 중이라면 run(ctx=ctx)는 그 실행을 재개하고 새 StartEvent를 보내지 않아요. 이전 실행이 완료되었다면 run(ctx=ctx)는 저장된 상태와 함께 새 실행을 시작합니다.

직렬화 가능한 상태

여기에 보관하는 상태는 실행을 내구성 있게(durable) 만들기 위해 스냅샷할 때 직렬화됩니다. 그러니 JSON 직렬화기로 인코딩할 수 있는 값으로 유지하세요. 클라이언트와 그 외 직렬화 불가능한 객체는 리소스에 두는 게 맞아요.

커스텀 값의 경우 Pydantic으로 직렬화 가능하게 만들거나, Context.to_dict()Context.from_dict()를 호출할 때 커스텀 직렬화기를 제공하면 됩니다.

더 알아보기