이벤트 스트림 처리
이벤트 스트림 처리 (Process Event Stream)
ProcessEventStream은 에이전트의 AgentStreamEvent 스트림 — 모델 스트리밍 및 도구 실행 이벤트 —을 핸들러로 전달하는 캐퍼빌리티예요. 등록하면 agent.run()이 자동으로 스트리밍을 활성화해서, 명시적인 event_stream_handler 인수 없이도 핸들러가 발동합니다.
출처: 문서
본문
실시간(realtime) 세션 중에는 스트림에 실시간 전용 RealtimeEvent 멤버도 포함됩니다.
from collections.abc import AsyncIterable
from pydantic_ai import Agent, AgentStreamEvent, RunContext
from pydantic_ai.capabilities import ProcessEventStream
async def log_events(ctx: RunContext, events: AsyncIterable[AgentStreamEvent]) -> None:
async for event in events:
print(event) # (1)
agent = Agent('openai:gpt-5.2', capabilities=[ProcessEventStream(log_events)])
- (1) 예를 들어 이벤트를 웹소켓, 진행 표시줄, 감사 로그(audit log)로 전달할 수 있어요.
핸들러는 두 가지 형태가 있습니다:
EventStreamHandler—None을 반환하는async def로, 위 예제와 같아요. 이벤트가 핸들러로 전달되며 그대로 통과하므로, 여러 핸들러(최상위event_stream_handler인수 포함)가 같은 스트림을 모두 관찰할 수 있습니다. 이벤트는 동기적으로 전달되므로 느린 핸들러는 스트림의 나머지에 백프레셔(back-pressure)를 걸어요.EventStreamProcessor— 이벤트를 생성하는 async generator예요. 생성된 것이 다운스트림 소비자에게 스트림을 대체하므로, 이벤트를 수정하거나 버리거나 추가할 수 있습니다.
캐퍼빌리티 등록은 다른 스트리밍 메커니즘과 함께 구성됩니다. Streaming all events에서 이벤트 어휘와 핸들러 예제를 확인하세요.
지속 실행 (Durable execution)
지속 실행 캐퍼빌리티 아래에서 ProcessEventStream 핸들러는 워크플로 코드에서 실행되며, 워크플로 재생(replay) 때 다시 실행되므로 결정적(deterministic) 이어야 해요. 도구 및 최종 출력 이벤트는 실시간으로 도착하고, 모델 이벤트는 각 모델 요청이 완료된 후 재생됩니다. 지속 경계 안에서 정확히 한 번만 실행되어야 하는 핸들러 I/O라면, 대신 지속성 캐퍼빌리티에 event_stream_handler=를 전달하세요.