이벤트 스트리밍

이벤트 스트리밍 (Event streaming)

메시지, 상태, 서브그래프, 출력, 확장 기능을 위한 타입 있는 프로젝션(typed projection)으로 LangGraph 런을 스트리밍해요.

이벤트 스트리밍은 대부분의 LangGraph 애플리케이션 코드에서 권장되는 프로세스 내(in-process) 스트리밍 모델이에요. 여러 방식으로 동시에 소비할 수 있는 런 스트림(run stream) 객체를 반환합니다.

출처: 문서

본문

빠른 시작 (Quickstart)

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를 참고하세요.

구성 요소가 어떻게 맞물리는지 (How the pieces fit together)

스트리밍 스택에는 두 개의 주요 레이어가 있어요:

  1. Streaming(저수준 스트리밍) — Pregel 엔진에서 원시 그래프 실행 이벤트를 방출합니다.
  2. Event streaming(이벤트 스트리밍) — 그 이벤트를 정규화하고 스트림 트랜스포머(stream transformer)를 통과시킨 뒤 타입 있는 프로젝션을 노출해요.
Pregel engine (Runs graph steps)
        ↓ emits
Raw Pregel events — updates, values, messages, custom, checkpoints, tasks, debug
        ↓ sent to
Event router (Routes each event through the transformer pipeline)
        ↓ cascades through
Stream transformers — ValuesTransformer, MessagesTransformer, ..., Custom transformers
        ↓ produces
Event Stream (Projected events for application code)

이벤트 라우터(event router)가 두 레이어 사이의 다리 역할을 해요. 정규화된 Pregel 이벤트를 받아 각 이벤트를 등록된 스트림 트랜스포머에 전달합니다. 내장 트랜스포머가 stream.messages, stream.values, stream.subgraphs, stream.output 같은 표준 프로젝션을 만들고, 커스텀 트랜스포머는 stream.extensions 아래에 애플리케이션 특화 프로젝션을 추가할 수 있어요.

이벤트 스트리밍이 제공하는 것 (What event streaming provides)

런 스트림은 하나의 기반 이벤트 흐름 위에 타입 있는 프로젝션을 노출합니다:

프로젝션 용도
stream 모든 프로토콜 이벤트를 순회.
stream.messages 채팅 모델 메시지와 토큰 델타 스트리밍.
stream.values 상태 스냅샷 순회 및 최종 값 대기.
stream.output 최종 출력 대기.
stream.subgraphs 중첩 그래프 실행 발견 및 관찰.
stream.interrupts human-in-the-loop 인터럽트 페이로드 점검.
stream.interrupted 런이 사람의 입력을 위해 멈췄는지 확인.
stream.extensions 커스텀 스트림 트랜스포머 프로젝션 소비.

여러 소비자가 이 프로젝션을 동시에 읽을 수 있어요. stream.messages를 읽는 것이 stream.values, stream.subgraphs, stream.output에 필요한 이벤트를 소모하지 않습니다.

이벤트 스트리밍은 streaming보다 한 단계 위에 있어요. 저수준 streaming은 updates, values, messages, custom, checkpoints, tasks, debug 같은 stream_mode를 통해 원시 그래프 실행 이벤트를 노출합니다. 그 모드들에 대한 저수준 접근이 필요하면 streaming을, 애플리케이션 코드가 타입 있는 프로젝션을 쓰면 이벤트 스트리밍을 쓰세요.

메시지 스트리밍 (Stream messages)

채팅 모델 출력에는 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는 동기 코드에서 순회 가능(iterable)해요. 토큰 단위 출력을 위해 순회하거나, 전체 텍스트를 위해 str(message.text)를 호출하면 됩니다.

message.reasoning은 추론 델타를, message.tool_calls는 도구 호출 인자 청크를 노출해요. 텍스트·추론·도구 호출 청크를 정확한 도착 순서로 필요하다면 각 프로젝션을 따로 순회하는 대신 메시지 스트림의 원시 이벤트를 순회하세요.

서브그래프 스트리밍 (Stream subgraphs)

중첩 그래프 작업을 네임스페이스 문자열 파싱 없이 관찰하려면 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 streaming(서브에이전트 스트림)과 LangChain agent streaming(도구 호출·미들웨어 이벤트)을 참고하세요.

상태 스트리밍 (Stream state)

각 단계 후 전체 상태 스냅샷을 스트리밍하려면 stream.values를 쓰세요:

stream = graph.stream_events(input, version="v3")

for snapshot in stream.values:
    print(snapshot)

final_state = stream.output

여러 프로젝션 스트리밍 (Stream multiple projections)

비동기 코드에서 동시 소비하려면 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}")

인터럽트 후 재개 (Resume after an interrupt)

그래프가 사람의 입력을 위해 멈추면 stream.interruptedstream.interrupts를 점검하고, Command와 함께 stream_events(..., version="v3")를 다시 호출해 재개하세요.

재개하려면 체크포인터로 컴파일된 그래프와 스레드 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 all protocol events)

원시 프로토콜 이벤트 스트림이 필요할 때는 런 객체 자체를 쓰세요:

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 봉투(envelope)예요. 트랜스포머의 process(event)가 받는 형태와 같습니다.

class ProtocolEvent(TypedDict):
    seq: int                    # strictly increasing within a run; use for ordering
    method: str                 # channel name: "messages", "values", "updates", "custom", "tools", "lifecycle", ...
    params: ProtocolEventParams


class ProtocolEventParams(TypedDict):
    namespace: list[str]        # path of "<name>:<runtime_id>" segments from the root graph; [] is the root
    timestamp: int              # wall-clock milliseconds; can drift, don't rely on for ordering
    data: Any                   # channel-specific payload; shape depends on `method`

namespace는 루트 그래프에서 이벤트를 방출한 범위까지의 경로예요. 루트는 빈 배열 []입니다. 각 하위 실행은 "name:runtime_id" 세그먼트를 하나 추가하므로, 서브그래프 안의 중첩 도구 호출은 ["researcher:6f4d", "tools:91ac"]처럼 보입니다. : 앞의 이름은 안정적인 그래프·노드 이름이고, 접미사는 호출별 런타임 ID입니다. 특정 하위 트리만 신경 쓸 때는 원시 이벤트를 직접 네임스페이스로 필터링하세요 — stream.subgraphs는 중첩 그래프 실행에 대해 이미 이 작업을 합니다.

채널과 이벤트 생명주기 (Channels and event lifecycle)

원시 이벤트는 채널(channel) 위에서 흐릅니다. 채널 이름이 이벤트의 method로 나타나고, 각 채널은 특정 이벤트 형태를 방출합니다.

채널 용도
values 전체 그래프 상태 스냅샷.
updates 노드별 상태 델타.
messages 콘텐츠 블록 중심의 채팅 모델 출력.
tools 도구 호출 시작·스트리밍 출력·종료·오류 이벤트.
lifecycle 런·서브그래프·서브에이전트 상태 변경.
checkpoints 분기와 time travel을 위한 가벼운 체크포인트 봉투.
input human-in-the-loop 입력 요청·응답.
tasks Pregel 태스크 생성·결과 이벤트.
custom 그래프 코드의 사용자 정의 페이로드.
custom:<name> 애플리케이션 정의 스트림 트랜스포머 출력.

타입 있는 프로젝션(stream.messages, stream.values 등)은 이 채널들에서 만들어집니다. 런 객체를 직접 순회하면 채널 이름이 원시 이벤트의 method 필드로 나타납니다.

메시지 (Messages)

messages 채널은 출력을 콘텐츠 블록(content block)으로 모델링합니다. 데이터의 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)

tools 채널은 도구 실행을 노출합니다. 데이터의 event 필드는 다음 중 하나입니다:

  • tool-started
  • tool-output-delta
  • tool-finished
  • tool-error

도구 이벤트는 도구 호출 ID로 연관되므로, 한 번의 도구 실행을 messages 채널의 원래 도구 호출 콘텐츠 블록에 다시 연결할 수 있어요.

생명주기 (Lifecycle)

lifecycle 채널은 루트 런, 서브그래프, 서브에이전트 상태를 추적합니다. 데이터의 event 필드는 다음 중 하나입니다:

  • started
  • running
  • completed
  • failed
  • interrupted

event 외에도 lifecycle 데이터에는 선택적인 graph_name, error, 그리고 하위 범위가 왜 시작됐는지(부모 도구 호출, 팬아웃 send, 엣지 전이)를 설명하는 cause가 포함될 수 있어요.

직접 프로젝션 만들기 (Build your own projection)

스트림 트랜스포머는 이벤트 스트리밍의 프로젝션 레이어예요. 프로토콜 이벤트를 관찰하고, 자체 상태를 유지하며, 런의 파생 뷰(도구 활동, 토큰 총합, 진행 이벤트, 아티팩트, 다른 프로토콜용 메시지 등)를 노출합니다. StreamChannel은 트랜스포머가 그 뷰를 발행하는 데 쓰는 프로젝션 원시형(primitive)입니다.

내장 프로젝션(stream.messages, stream.values, stream.subgraphs, stream.output)과 제품별 프로젝션(LangChain의 stream.tool_calls, Deep Agents의 stream.subagents)도 같은 계약을 쓰는 트랜스포머예요. 사용자 트랜스포머는 컴파일 시점 또는 호출 시점 등록을 통해 위에 쌓이고, 그 프로젝션은 stream.extensions 아래에 나타납니다.

기존 프로젝션이 애플리케이션이 필요한 형태와 맞지 않을 때 직접 작성하세요.

트랜스포머가 어떻게 동작하는가 (How transformers work)

이벤트 스트리밍은 LangGraph Pregel 엔진의 스트리밍 출력으로 시작해요. 런타임이 그 청크를 프로토콜 이벤트로 정규화하고, 스트림 핸들러가 각 이벤트를 스트림 트랜스포머 스택을 통해 라우팅합니다.

Pregel modes --> Events --> Built-in projections --> User transformers --> Run projections

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

  1. 등록된 각 트랜스포머의 process(event) 훅을 순서대로 호출.
  2. 이름 있는 StreamChannel push를 프로토콜 이벤트 스트림으로 다시 연결.
  3. 트랜스포머가 억제하지 않는 한 이벤트를 런 스트림에 저장.
  4. 런이 끝나면 모든 트랜스포머의 finalize() 또는 fail() 호출.

트랜스포머는 관찰용(observational)이에요. 그래프 런타임을 다시 호출하지 않습니다. 대신 이벤트를 소비하고 파생 값을 StreamChannel, promise, 또는 다른 프로젝션 객체로 밀어 넣습니다.

트랜스포머 형태 (Transformer shape)

트랜스포머는 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()은 오류를 비채널 프로젝션에 전파합니다.

필수 스트림 모드 선언 (Declaring required stream modes)

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()
각 push도 메인 이벤트 스트림으로 흘려보내기 StreamChannel(name)

이름 있는 채널 페이로드는 직렬화 가능해야 해요. push된 각 값이 메인 스트림에서 custom:<name> 프로토콜 이벤트가 되기 때문입니다. promise, 비동기 iterable, 클래스 인스턴스, 그 밖의 프로세스 내 핸들은 이름 없는 채널에 두세요.

스트림 핸들러가 채널 생명주기를 소유합니다. init()이 채널을 반환하면 런이 끝날 때 핸들러가 그 채널을 닫거나 실패 처리해 줘요. 트랜스포머는 값만 push하면 됩니다.

예시: 이름 있는 채널 (Example: named channel)

StreamChannel에 문자열 이름을 전달하면 stream.extensions를 통해 스트리밍 프로젝션을 노출하는 동시에, push된 각 값을 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

예시: 이름 없는 채널 (Example: unnamed channel)

이름이 없으면 채널은 사이드 채널 프로젝션일 뿐입니다 — stream.extensions에서 접근할 수 있지만 원시 이벤트를 순회하는 소비자에게는 보이지 않아요. 메인 이벤트 스트림에 직렬화할 수 없는 프로세스 내 핸들(promise, 비동기 iterable, 클래스 인스턴스)을 담는 프로젝션에 적합합니다.

아래 예시는 이름 없는 채널을 get_stream_writer와 짝지어서, 그래프 노드가 방출하는 custom 채널 이벤트를 트랜스포머가 프로젝션으로 배출(drain)하게 합니다:

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)

예시: 최종값 프로젝션 (Example: final-value projection)

프로젝션이 메인 이벤트 스트림으로 흐르지 않아야 할 때는 이름 없는 스트림, promise, 그 밖의 프로세스 내 객체를 쓰세요:

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()

호출 시점 또는 컴파일 시점 등록 (Register at call time or compile time)

로컬 실험을 위해 호출 시점에 트랜스포머를 전달하세요:

stream = graph.stream_events(
    input,
    version="v3",
    transformers=[StatsTransformer, ToolActivityTransformer],
)

그 그래프의 모든 런이 프로젝션을 만들어야 한다면 트랜스포머를 컴파일 시점에 넣으세요:

graph = builder.compile(
    transformers=[StatsTransformer, ToolActivityTransformer],
)

내장: ToolCallTransformer (Built-in: 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)