이벤트 스트리밍 — 진행 상황을 실시간으로 전달하기
이벤트 스트리밍 — 진행 상황을 실시간으로 전달하기
LLM이 답을 생성하는 동안 사용자에게 "지금 이러고 있어요" 같은 중간 상황을 보여주고 싶다면, 워크플로 안에서 이벤트 스트림으로 진행 상황을 흘려보내면 돼요. 스텝 안에서는 ctx.write_event_to_stream(...), 워크플로 바깥에서는 handler.stream_events()를 쓰면 됩니다.
출처: 공식문서
스트림에 이벤트 쓰기
workflow.run(...)은 워크플로를 시작하고 WorkflowHandler를 돌려줘요. 이 핸들은 최종 결과를 위해 await할 수 있고, 해당 실행의 이벤트 스트림도 소유하고 있어요.
간단한 3스텝 워크플로를 예로 들어볼게요. 진행 상황을 담을 ProgressEvent를 하나 정의해요.
from llama_index.core.workflow import Workflow, Context, step
from llama_index.core.workflow.events import Event, StartEvent, StopEvent
class FirstEvent(Event):
first_output: str
class SecondEvent(Event):
second_output: str
response: str
class ProgressEvent(Event):
msg: str
1번 스텝과 3번 스텝에서는 개별 이벤트를 스트림에 직접 써요. 2번 스텝은 LLM의 astream_complete로 생성되는 응답 청크를 하나씩 받으면서, 그 조각마다 ctx.write_event_to_stream(ProgressEvent(msg=response.delta))로 스트림에 흘려보내요. 이렇게 하면 최종 응답이 완성되기 전에 토큰 조각들이 실시간으로 출력돼요.
class MyWorkflow(Workflow):
@step
async def step_one(self, ctx: Context, ev: StartEvent) -> FirstEvent:
ctx.write_event_to_stream(ProgressEvent(msg="Step one is happening"))
return FirstEvent(first_output="First step complete.")
OpenAI()는 환경에 OPENAI_API_KEY가 설정돼 있다고 가정해요. 없으면 api_key 파라미터로 넘길 수도 있어요.
바깥에서 스트림 읽기
스트림을 읽는 쪽은 handler.stream_events()로 이벤트를 차례로 순회해요. 이 스트림은 StopEvent가 전달될 때까지 쓰여진 모든 이벤트를 내보내고, 그 뒤에 핸들을 await하면 최종 결과를 얻어요.
async def main():
w = MyWorkflow(timeout=30, verbose=True)
handler = w.run(first_input="Start the workflow.")
async for ev in handler.stream_events():
if isinstance(ev, ProgressEvent):
print(ev.msg)
final_result = await handler
워크플로 종료 상황 처리하기
스트림은 진행 상황뿐 아니라 워크플로의 종료 상황도 이벤트로 알려줘요.
WorkflowTimedOutEvent— 워크플로가 타임아웃을 넘겼을 때 게시돼요.timeout(초)과 그때 실행 중이던 스텝 이름 목록active_steps를 담아요.WorkflowCancelledEvent— 사용자가 워크플로를 취소했을 때 게시돼요.WorkflowFailedEvent— 재시도를 다 소진한 뒤 스텝이 영구 실패했을 때 게시돼요.step_name,exception,attempts,elapsed_seconds를 담아요.
from llama_index.core.workflow.events import (
WorkflowTimedOutEvent, WorkflowCancelledEvent, WorkflowFailedEvent,
)
async for ev in handler.stream_events():
if isinstance(ev, WorkflowTimedOutEvent):
print(f"Workflow timed out after {ev.timeout}s")
elif isinstance(ev, WorkflowCancelledEvent):
print("Workflow was cancelled")
elif isinstance(ev, WorkflowFailedEvent):
print(f"Step '{ev.step_name}' failed after {ev.attempts} attempts: {ev.exception}")
더 알아보기
- 워크플로 동시 실행 — 병렬 처리
- 에러 처리(재시도) — 실패와 재시도 다루기
- 비동기 워크플로 작성 — async 스텝 이해하기