이벤트 스트리밍

이벤트 스트리밍

메시지, 상태, 서브그래프, 출력, 확장을 위한 타입 프로젝션으로 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를 보세요.

출처: 이벤트 스트리밍 - 공식 문서

구성 요소가 어떻게 맞물리는지

스트리밍 스택은 두 개의 메인 레이어로 나뉘어요:

  1. Streaming은 Pregel 엔진에서 원시 그래프 실행 이벤트를 내보내요.
  2. 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.interruptedstream.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-start
  • content-block-start
  • content-block-delta
  • content-block-finish
  • message-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-started
  • tool-output-delta
  • tool-finished
  • tool-error

툴 이벤트는 툴 콜 ID로 상관관계가 있어, 툴 실행을 messages 채널의 원래 툴콜 콘텐츠 블록으로 이어붙일 수 있어요.

라이프사이클

lifecycle 채널은 루트 실행·서브그래프·서브에이전트 상태를 추적해요. 데이터의 event 필드는 다음 중 하나예요:

  • started
  • running
  • completed
  • failed
  • interrupted

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]

스트림 핸들러는 한 스트림의 중앙 디스패처예요. 모든 프로토콜 이벤트에 대해 다음을 해요:

  1. 등록된 각 트랜스포머의 process(event) 훅을 순서대로 호출해요.
  2. 이름 있는 StreamChannel 푸시를 프로토콜 이벤트 스트림에 다시 연결해요.
  3. 트랜스포머가 이를 억제하지 않으면 이벤트를 실행 스트림에 저장해요.
  4. 실행이 끝나면 모든 트랜스포머의 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에서 스트리밍을 쓸 때는 관련 제품 문서를 보세요:

와이어 레벨 이벤트·명령 형식은 Agent Protocol 저장소에 정의되어 있고, PyPI의 langchain-protocol과 npm의 @langchain/protocol로 소비할 수 있어요.

더 알아보기 (Learn more)

  • Streaming — 저수준 stream_mode 기반 스트리밍
  • Persistence — 재개와 인터럽트를 위한 체크포인팅