이벤트 스트리밍
이벤트 스트리밍
메시지, 상태, 서브그래프, 출력, 확장을 위한 타입 프로젝션으로 LangGraph 실행을 스트리밍해요.
이벤트 스트리밍은 대부분의 LangGraph 애플리케이션 코드에서 권장하는 인프로세스 스트리밍 모델이에요. 실행 스트림 객체를 반환하는데, 이 객체를 여러 방식으로 동시에 소비할 수 있죠.
빠른 시작
stream = graph.stream_events({
"messages": [{"role": "user", "content": "What is 42 * 17?"}],
}, version="v3")
for message in stream.messages:
for token in message.text:
print(token, end="", flush=True)
final_state = stream.output
Agent Server 뒤에 배포된 그래프를 대상으로 스트리밍하려면 LangSmith Streaming API를 보세요.
출처: 이벤트 스트리밍 - 공식 문서
구성 요소가 어떻게 맞물리는지
스트리밍 스택은 두 개의 메인 레이어로 나뉘어요:
- Streaming은 Pregel 엔진에서 원시 그래프 실행 이벤트를 내보내요.
- Event streaming은 그 이벤트를 정규화하고, 스트림 트랜스포머를 통과시킨 다음, 타입 프로젝션으로 노출해요.
Pregel 엔진 (그래프 단계 실행)
│ emits
▼
원시 Pregel 이벤트 (updates, values, messages, custom, checkpoints, tasks, debug)
│ sent to
▼
이벤트 라우터 (각 이벤트를 트랜스포머 파이프라인으로 라우팅)
│ cascades through
▼
스트림 트랜스포머 (ValuesTransformer, MessagesTransformer, ..., Custom transformers)
│ produces
▼
이벤트 스트림 (애플리케이션 코드를 위한 프로젝트 이벤트)
이벤트 라우터는 두 레이어 사이의 다리예요. 정규화된 Pregel 이벤트를 받아 각 이벤트를 등록된 스트림 트랜스포머에 통과시키죠. 내장 트랜스포머는 stream.messages, stream.values, stream.subgraphs, stream.output 같은 표준 프로젝션을 만들어요. 커스텀 트랜스포머는 stream.extensions 아래에 애플리케이션 특화 프로젝션을 추가할 수 있어요.
이벤트 스트리밍이 제공하는 것
실행 스트림은 하나의 기본 이벤트 흐름 위에 타입 프로젝션을 노출해요:
| 프로젝션 | 용도 |
|---|---|
stream |
모든 프로토콜 이벤트를 순회해요. |
stream.messages |
채팅 모델 메시지와 토큰 델타를 스트리밍해요. |
stream.values |
상태 스냅샷을 순회하고 최종 값을 기다려요. |
stream.output |
최종 출력을 기다려요. |
stream.subgraphs |
중첩 그래프 실행을 발견하고 관찰해요. |
stream.interrupts |
휴먼 인 더 루프 인터럽트 페이로드를 검사해요. |
stream.interrupted |
실행이 인간 입력을 위해 멈췄는지 확인해요. |
stream.extensions |
커스텀 스트림 트랜스포머 프로젝션을 소비해요. |
여러 소비자가 이 프로젝션을 동시에 읽을 수 있어요. stream.messages를 읽어도 stream.values, stream.subgraphs, stream.output이 필요한 이벤트를 소비하지 않아요.
이벤트 스트리밍은 원시 그래프 실행 이벤트를 stream_mode(updates, values, messages, custom, checkpoints, tasks, debug 등)로 노출하는 streaming보다 한 단계 위에 있어요. 저수준 모드 접근이 필요하면 streaming을, 애플리케이션 코드가 타입 프로젝션을 쓰면 이벤트 스트리밍을 쓰세요.
메시지 스트리밍
채팅 모델 출력에는 stream.messages를 써요:
stream = graph.stream_events(input, version="v3")
for message in stream.messages:
text = str(message.text)
usage = message.output.usage_metadata
print(text)
print(usage)
message.text는 동기 코드에서 순회 가능해요. 토큰 단위 출력으로 순회하거나 str(message.text)로 전체 텍스트를 얻어요.
message.reasoning은 추론 델타를, message.tool_calls는 툴 콜 인수 청크를 노출해요. 텍스트·추론·툴콜 청크를 정확한 도착 순서대로 원한다면, 각 프로젝션을 따로 순회하지 말고 메시지 스트림의 원시 이벤트를 순회하세요.
서브그래프 스트리밍
네임스페이스 문자열을 파싱하지 않고 중첩 그래프 작업을 관찰하려면 stream.subgraphs를 써요:
stream = graph.stream_events(input, version="v3")
for subgraph in stream.subgraphs:
print(subgraph.graph_name, subgraph.path)
for message in subgraph.messages:
print(message.text)
subgraph.graph_name은 컴파일된 그래프나 에이전트의 name이에요. 툴에서 발송된 이름 있는 에이전트(예: Deep Agents의 task 툴로 호출된 create_agent(name=...))는 그 이름으로 여기에 나타나고, 스코프를 여는 lifecycle 이벤트는 발송한 툴 콜을 가리키는 cause를 가져요. Lifecycle을 참고하세요.
제품별 스트림은 Deep Agents 스트리밍(서브에이전트 스트림)과 LangChain 에이전트 스트리밍(툴 콜·미들웨어 이벤트)을 보세요.
상태 스트리밍
각 단계 후 전체 상태 스냅샷을 스트리밍하려면 stream.values를 써요:
stream = graph.stream_events(input, version="v3")
for snapshot in stream.values:
print(snapshot)
final_state = stream.output
여러 프로젝션 스트리밍
비동기 코드에서 동시 소비하려면 asyncio.gather와 함께 astream_events를 써요:
import asyncio
stream = await graph.astream_events(input, version="v3")
async def consume_messages():
async for message in stream.messages:
print(f"[llm] node={message.node}")
async def consume_subgraphs():
async for subgraph in stream.subgraphs:
print(f"[subgraph] path={subgraph.path}")
await asyncio.gather(consume_messages(), consume_subgraphs())
동기 코드에서는 stream.interleave(...)로 여러 프로젝션을 엄격한 도착 순서로 소비해요:
stream = graph.stream_events(input, version="v3")
for name, item in stream.interleave("values", "messages", "subgraphs"):
if name == "values":
print(f"[state] keys={list(item)}")
elif name == "messages":
print(f"[llm] node={item.node}")
elif name == "subgraphs":
print(f"[subgraph] path={item.path}")
인터럽트 후 재개
그래프가 인간 입력을 위해 멈추면 stream.interrupted와 stream.interrupts를 검사하고, Command를 들고 다시 stream_events(..., version="v3")를 호출해 재개해요.
재개하려면 체크포인터로 컴파일된 그래프와 thread ID를 담은 config가 필요해요. persistence를 보세요.
from langgraph.types import Command
stream = graph.stream_events(input, version="v3")
for message in stream.messages:
print(message.text)
if stream.interrupted:
print(stream.interrupts)
stream = graph.stream_events(
Command(resume={"decisions": [{"type": "approve"}]}),
version="v3",
)
final_state = stream.output
모든 프로토콜 이벤트 스트리밍
원시 프로토콜 이벤트 스트림을 원하면 실행 객체 자체를 써요:
stream = graph.stream_events({
"messages": [{"role": "user", "content": "What is 42 * 17?"}],
}, version="v3")
for event in stream:
namespace = event["params"]["namespace"]
print(namespace, event["method"], event["params"]["data"])
각 이벤트는 채널별 페이로드를 감싸는 ProtocolEvent 엔벨로프예요. 같은 모양이 트랜스포머의 process(event)가 받는 것이기도 해요.
class ProtocolEvent(TypedDict):
seq: int # 실행 내에서 엄격히 증가; 정렬에 사용
method: str # 채널 이름: "messages", "values", "updates", "custom", "tools", "lifecycle", ...
params: ProtocolEventParams
class ProtocolEventParams(TypedDict):
namespace: list[str] # 루트 그래프부터의 "<name>:<runtime_id>" 세그먼트 경로; []는 루트
timestamp: int # 벽시계 밀리초; 드리프트할 수 있으니 정렬에 의존하지 말 것
data: Any # 채널별 페이로드; 모양은 `method`에 따라 다름
namespace는 루트 그래프에서 이벤트를 발생시킨 스코프까지의 경로예요. 루트는 빈 배열 []이죠. 각 하위 실행이 "name:runtime_id" 세그먼트를 하나 추가하므로, 서브그래프 안 중첩 툴 콜은 ["researcher:6f4d", "tools:91ac"]처럼 보여요. : 앞의 이름은 안정적인 그래프 또는 노드 이름이고, 접미사는 호출별 런타임 ID예요. 특정 서브트리에만 관심이 있으면 네임스페이스로 원시 이벤트를 직접 필터링하세요. stream.subgraphs는 중첩 그래프 실행에 대해 이미 이 작업을 해줘요.
채널과 이벤트 라이프사이클
원시 이벤트는 채널로 흘러요. 채널 이름이 이벤트의 method로 나타나고, 각 채널은 특정 이벤트 모양을 내보내요.
| 채널 | 용도 |
|---|---|
values |
전체 그래프 상태 스냅샷. |
updates |
노드별 상태 델타. |
messages |
콘텐츠 블록 중심의 채팅 모델 출력. |
tools |
툴 콜 시작, 스트리밍 출력, 종료, 오류 이벤트. |
lifecycle |
실행·서브그래프·서브에이전트 상태 변경. |
checkpoints |
브랜칭·타임 트래블용 경량 체크포인트 엔벨로프. |
input |
휴먼 인 더 루프 입력 요청과 응답. |
tasks |
Pregel 태스크 생성과 결과 이벤트. |
custom |
그래프 코드에서 온 사용자 정의 페이로드. |
custom:<name> |
애플리케이션 정의 스트림 트랜스포머 출력. |
타입 프로젝션(stream.messages, stream.values 등)은 이 채널들로 만들어져요. 실행 객체를 직접 순회할 때 채널 이름이 원시 이벤트의 method 필드로 나타나요.
메시지
messages 채널은 출력을 콘텐츠 블록으로 모델링해요. 데이터의 event 필드는 다음 중 하나예요:
message-startcontent-block-startcontent-block-deltacontent-block-finishmessage-finish
콘텐츠 블록은 명시적인 경계를 가져요: 블록이 시작하고, 0개 이상의 델타를 내고, 같은 메시지의 다음 블록이 시작하기 전에 끝나죠. 이 덕분에 토큰 스트리밍, 추론 블록, 툴콜 블록, 멀티모달 콘텐츠가 프로바이더별 형식 없이 명시적으로 표현돼요. message-finish는 토큰 사용량을 포함할 수 있고, 복구 불가능한 모델 콜 실패는 메시지 오류 이벤트로 도착해요.
stream.messages 프로젝션 대신 원시 콘텐츠 블록 이벤트를 직접 소비하려면:
for event in stream:
if event["method"] != "messages":
continue
data = event["params"]["data"][0]
if not isinstance(data, dict):
continue
if data.get("event") != "content-block-delta":
continue
block = data.get("delta") or {}
if block.get("type") == "text-delta":
print(block.get("text", ""), end="", flush=True)
elif block.get("type") == "reasoning-delta":
print(f"[thinking]{block.get('reasoning', '')}", end="", flush=True)
툴
tools 채널은 툴 실행을 노출해요. 데이터의 event 필드는 다음 중 하나예요:
tool-startedtool-output-deltatool-finishedtool-error
툴 이벤트는 툴 콜 ID로 상관관계가 있어, 툴 실행을 messages 채널의 원래 툴콜 콘텐츠 블록으로 이어붙일 수 있어요.
라이프사이클
lifecycle 채널은 루트 실행·서브그래프·서브에이전트 상태를 추적해요. 데이터의 event 필드는 다음 중 하나예요:
startedrunningcompletedfailedinterrupted
event 외에 라이프사이클 데이터는 선택적 graph_name, error, 그리고 하위 스코프가 왜 시작됐는지(부모 툴 콜, 팬아웃 send, 엣지 전이)를 설명하는 cause를 포함할 수 있어요.
나만의 프로젝션 만들기
스트림 트랜스포머는 이벤트 스트리밍의 프로젝션 레이어예요. 프로토콜 이벤트를 관찰하고, 자신의 상태를 유지하며, 실행의 파생 뷰(툴 활동, 토큰 합계, 진행 이벤트, 산출물, 다른 프로토콜용 메시지 등)를 노출하죠. 트랜스포머가 그 뷰를 게시하는 데 쓰는 프로젝션 기본 요소가 StreamChannel이에요.
내장 프로젝션(stream.messages, stream.values, stream.subgraphs, stream.output)과 제품별 프로젝션(LangChain의 stream.tool_calls, Deep Agents의 stream.subagents)은 같은 계약을 쓰는 트랜스포머예요. 사용자 트랜스포머는 컴파일 시점 또는 호출 시점 등록으로 그 위에 쌓이고, 프로젝션은 stream.extensions 아래에 나타나요.
기존 프로젝션이 애플리케이션이 필요로 하는 모양과 맞지 않을 때 하나 작성하세요.
트랜스포머가 어떻게 동작하는지
이벤트 스트리밍은 LangGraph Pregel 엔진의 스트리밍 출력에서 시작해요. 런타임이 그 청크를 프로토콜 이벤트로 정규화하고, 스트림 핸들러가 각 이벤트를 스트림 트랜스포머 스택에 라우팅해요.
flowchart TD
A[Pregel modes] --> B[Events]
B --> C[Built-in projections]
C --> D[User transformers]
D --> E[Run projections]
스트림 핸들러는 한 스트림의 중앙 디스패처예요. 모든 프로토콜 이벤트에 대해 다음을 해요:
- 등록된 각 트랜스포머의
process(event)훅을 순서대로 호출해요. - 이름 있는
StreamChannel푸시를 프로토콜 이벤트 스트림에 다시 연결해요. - 트랜스포머가 이를 억제하지 않으면 이벤트를 실행 스트림에 저장해요.
- 실행이 끝나면 모든 트랜스포머의
finalize()또는fail()을 호출해요.
트랜스포머는 관찰적이에요. 그래프 런타임에 다시 콜백하지 않고, 이벤트를 소비해 파생 값을 StreamChannel, 프로미스, 또는 다른 프로젝션 객체에 푸시해요.
트랜스포머 모양
트랜스포머는 StreamTransformer 인터페이스를 구현해요:
from langgraph.stream import ProtocolEvent, StreamTransformer
class MyTransformer(StreamTransformer):
def init(self) -> dict:
...
def process(self, event: ProtocolEvent) -> bool:
...
def finalize(self) -> None:
...
def fail(self, err: BaseException) -> None:
...
init()는 프로젝션 객체를 만들어요. 사용자 트랜스포머 프로젝션은stream.extensions아래에 나타나요.process()는 각 프로토콜 이벤트를 관찰해요.ProtocolEvent모양은 모든 프로토콜 이벤트 스트리밍을 보세요. 의도적으로 원본 이벤트를 억제하려 할 때만false를 반환해요.finalize()는 성공적인 스트림 후 채널이 아닌 프로젝션을 닫거나 해소해요.fail()은 오류를 채널이 아닌 프로젝션에 전파해요.
필요한 스트림 모드 선언
required_stream_modes는 기본 그래프가 스트림 동안 내보낼 Pregel 스트림 모드를 제어해요. 런타임은 등록된 모든 트랜스포머의 required_stream_modes의 합집합을 구해, 그 합집합을 그래프 .stream() 호출의 stream_mode 인수로 넘겨요. 어떤 트랜스포머도 요청하지 않는 모드는 절대 내보내지 않아요 — ("custom",)을 선언해야 custom 이벤트가 실행을 통해 흐르는 거예요.
class CustomTransformer(StreamTransformer):
required_stream_modes = ("custom",) # [!code highlight]
def process(self, event: ProtocolEvent) -> bool:
if event["method"] == "custom":
...
return True
process()는 그래프가 내보내는 모든 이벤트를 받고 event["method"]로 필터링하는 책임이 있어요. 선언은 업스트림 방출을 켜는 것이지, process()가 보는 범위를 좁히는 게 아니에요. 유효한 값은 Pregel 스트림 모드인 "messages", "tools", "custom", "values", "updates", "checkpoints", "tasks", "debug"예요. 각 트랜스포머는 자신이 다루는 모든 모드를 선언해야 해요. 빠뜨린 모드는 그래프가 내보내지 않아 process()에 도달하지 못해요.
StreamChannel
StreamChannel은 트랜스포머가 값을 스트리밍하는 데 쓰는 프로젝션 기본 요소예요. 항상 stream.extensions.<name>에 순회 가능한 스트림을 노출해요. 생성자 인자가 각 push()가 실행의 메인 이벤트 스트림에도 custom:<name> 이벤트로 흐르는지 — 즉 프로젝션 값이 원시 프로토콜 이벤트를 순회할 때 보이는지 —를 결정해요.
| 필요 | 사용 |
|---|---|
| 사이드 채널 프로젝션만 | StreamChannel() |
| 각 푸시를 메인 이벤트 스트림에도 | StreamChannel(name) |
이름 있는 채널 페이로드는 직렬화 가능해야 해요. 각 푸시된 값이 메인 스트림의 custom:<name> 프로토콜 이벤트가 되기 때문이죠. 프로미스, 비동기 순회자, 클래스 인스턴스 등 인프로세스 핸들은 이름 없는 채널에 두세요.
스트림 핸들러가 채널 라이프사이클을 소유해요. init()이 채널을 반환하면, 실행이 끝날 때 핸들러가 닫거나 실패시켜줘요. 트랜스포머는 값만 푸시해요.
예제: 이름 있는 채널
StreamChannel에 문자열 이름을 넘기면 stream.extensions을 통해 스트리밍 프로젝션을 노출하고, 각 푸시된 값을 실행의 메인 이벤트 스트림에 custom:<name> 프로토콜 이벤트로 전달해요:
from typing import TypedDict
from langgraph.stream import ProtocolEvent, StreamChannel, StreamTransformer
class ToolActivity(TypedDict):
name: str
status: str
class ToolActivityTransformer(StreamTransformer):
required_stream_modes = ("tools",)
def __init__(self, scope: tuple[str, ...] = ()) -> None:
super().__init__(scope)
self.activity = StreamChannel[ToolActivity]("tool_activity")
def init(self) -> dict:
return {"tool_activity": self.activity}
def process(self, event: ProtocolEvent) -> bool:
if event["method"] != "tools":
return True
data = event["params"]["data"]
if isinstance(data, dict) and data.get("tool_name") and data.get("event"):
status = "error" if data["event"] == "tool-error" else "started"
self.activity.push({"name": data["tool_name"], "status": status})
return True
예제: 이름 없는 채널
이름이 없으면 채널은 사이드 채널 프로젝션만 돼요. stream.extensions에서 접근 가능하지만 원시 이벤트를 순회하는 소비자에게는 보이지 않죠. 메인 이벤트 스트림에 직렬화할 수 없는 인프로세스 핸들(프로미스, 비동기 순회자, 클래스 인스턴스)을 담는 프로젝션에 맞는 선택이에요.
아래 예제는 이름 없는 채널을 get_stream_writer와 짝지어, 그래프 노드가 custom-채널 이벤트를 내보내면 트랜스포머가 그걸 프로젝션으로 배출하도록 해요:
from langgraph.config import get_stream_writer
from langgraph.stream import ProtocolEvent, StreamChannel, StreamTransformer
def node(state):
writer = get_stream_writer()
writer({"kind": "progress", "message": "retrieving context"})
return state
class CustomTransformer(StreamTransformer):
required_stream_modes = ("custom",)
def __init__(self, scope: tuple[str, ...] = ()) -> None:
super().__init__(scope)
self.log = StreamChannel()
def init(self) -> dict:
return {"custom": self.log}
def process(self, event: ProtocolEvent) -> bool:
if event["method"] == "custom":
self.log.push(event["params"]["data"])
return True
stream = graph.stream_events(input, version="v3", transformers=[CustomTransformer])
for item in stream.extensions["custom"]:
print(item)
예제: 최종 값 프로젝션
프로젝션이 메인 이벤트 스트림으로 흐르지 않아야 할 때는 이름 없는 스트림, 프로미스, 또는 다른 인프로세스 객체를 써요:
from langgraph.stream import ProtocolEvent, StreamChannel, StreamTransformer
class StatsTransformer(StreamTransformer):
required_stream_modes = ("messages",)
def __init__(self, scope: tuple[str, ...] = ()) -> None:
super().__init__(scope)
self.total_tokens = 0
self.total_tokens_log = StreamChannel[int]()
def init(self) -> dict:
return {"total_tokens": self.total_tokens_log}
def process(self, event: ProtocolEvent) -> bool:
data = event["params"]["data"]
if isinstance(data, dict):
usage = data.get("usage") or {}
self.total_tokens += usage.get("output_tokens") or 0
return True
def finalize(self) -> None:
self.total_tokens_log.push(self.total_tokens)
self.total_tokens_log.close()
호출 시점 또는 컴파일 시점 등록
로컬 실험용으로는 호출 시점에 트랜스포머를 넘겨요:
stream = graph.stream_events(
input,
version="v3",
transformers=[StatsTransformer, ToolActivityTransformer],
)
그래프의 모든 실행이 그 프로젝션을 만들어야 하면 컴파일 시점에 트랜스포머를 넣어요:
graph = builder.compile(
transformers=[StatsTransformer, ToolActivityTransformer],
)
내장: ToolCallTransformer
LangGraph는 ToolCallTransformer를 내장으로 제공해요. 등록하면 일반 StateGraph에서 stream.tool_calls를 노출해요:
from langgraph.prebuilt import ToolCallTransformer
stream = graph.stream_events(input, version="v3", transformers=[ToolCallTransformer])
for tool_call in stream.tool_calls:
print(tool_call.tool_name, tool_call.input)
관련 문서
LangGraph는 스트리밍 기본 요소를 정의해요. LangChain이나 Deep Agents에서 스트리밍을 쓸 때는 관련 제품 문서를 보세요:
- LangChain 에이전트 스트리밍은 ReAct 스타일 에이전트 메시지, 툴 콜, 미들웨어 갱신을 다뤄요.
- Deep Agents 스트리밍은 서브에이전트, 중첩 메시지, 서브에이전트 툴 콜을 다뤄요.
- LangChain 프론트엔드 패턴과 LangGraph 프론트엔드 패턴은 스트리밍된 상태 위에 만든 UI 사용 사례를 보여줘요.
- LangSmith Streaming API는 Agent Server 뒤에 배포된 그래프에 대한 스트리밍을 다뤄요.
와이어 레벨 이벤트·명령 형식은 Agent Protocol 저장소에 정의되어 있고, PyPI의 langchain-protocol과 npm의 @langchain/protocol로 소비할 수 있어요.
더 알아보기 (Learn more)
- Streaming — 저수준
stream_mode기반 스트리밍 - Persistence — 재개와 인터럽트를 위한 체크포인팅