이벤트 스트리밍 API: 배포된 LangGraph를 타이핑된 투영으로 쏟아내는 법

이벤트 스트리밍 API: 배포된 LangGraph를 타이핑된 투영으로 쏟아내는 법

LangSmith Deployment의 이벤트 스트리밍은 타이핑된 투영(typed projection) 방식의 스트리밍 모델입니다. LangGraph SDK(Python/JavaScript)가 LangSmith Deployment API에 단일 구독을 열어 두고, 한 run에서 동시에 소비할 수 있는 타이핑된 투영들—메시지, 상태, 도구 호출, 서브그래프, 출력, 커스텀 트랜스포머 확장—을 노출합니다. 이벤트 스트리밍은 원시 스트림 모드를 노출하는 streaming API보다 한 단계 위에 있습니다.

출처: Event streaming API 공식 문서

이벤트 스트리밍을 쓰려면 LangGraph Agent Server가 langgraph-api>=0.10.0이어야 합니다. 관리형 LangSmith 배포는 자동 갱신되지만, 셀프호스트 서버는 호환 버전이어야 합니다. 클라이언트 SDK는 langgraph-sdk>=0.4.0 + langchain-core>=1.4.0(Python) 또는 @langchain/langgraph-sdk>=1.9.15(JavaScript)여야 합니다. 이전 버전 서버는 계속 기존 streaming API를 서빙합니다.

퀵스타트

Python:

from langgraph_sdk import get_client

client = get_client(url=DEPLOYMENT_URL, api_key=API_KEY)

async with client.threads.stream(assistant_id="agent") as thread:
    await thread.run.start(
        input={"messages": [{"role": "user", "content": "What is 42 * 17?"}]},
    )

    async for message in thread.messages:
        async for token in message.text:
            print(token, end="", flush=True)

    final_state = await thread.output

JavaScript:

import { Client } from "@langchain/langgraph-sdk";

const client = new Client({
  apiUrl: process.env.DEPLOYMENT_URL,
  apiKey: process.env.API_KEY,
});

const thread = client.threads.stream({ assistantId: "agent" });

await thread.run.start({
  input: { messages: [{ role: "user", content: "What is 42 * 17?" }] },
});

for await (const message of thread.messages) {
  for await (const token of message.text) {
    process.stdout.write(token);
  }
}

const finalState = await thread.output;
await thread.close();

JavaScript 스트림에는 async with에 해당하는 게 없으므로, 쓰고 나면 await thread.close()로 기저 구독을 해제하세요.

기본적으로 SDK는 Server-Sent Events로 스트리밍합니다. 대신 양방향 WebSocket을 쓰려면 client.threads.stream(...)transport="websocket"를 넘기세요.

이벤트 스트리밍이 제공하는 것

client.threads.stream(...)이 돌려주는 스트림은 하나의 기저 이벤트 흐름 위에 여러 타이핑된 투영을 노출합니다.

투영 용도
thread.events 모든 원시 프로토콜 이벤트 순회(Python). JavaScript는 thread.subscribe(...)를 엶.
thread.messages 채팅 모델 메시지, 토큰 델타, 추론, 도구 호출 인자 청크 스트리밍.
thread.values 상태 스냅샷 순회 및 최종 값 대기.
thread.output 최종 출력 대기.
thread.tool_calls(thread.toolCalls) 조립된 입력·스트리밍 출력·결과를 가진 도구 호출 관찰.
thread.subgraphs 중첩 그래프 실행 발견·관찰.
thread.subagents thread.subgraphs의 서브에이전트 관점. Deep Agents 서브에이전트 호출 처리 시 사용.
thread.interrupts 인간-루프(human-in-the-loop) 인터럽트 페이로드 검사.
thread.interrupted run이 인간 입력을 위해 일시정지했는지 확인.
thread.extensions custom:<name> 채널에 발행된 커스텀 스트림 트랜스포머 투영 소비.

여러 소비자가 이 투영들을 동시에 읽을 수 있습니다. thread.messages를 읽는다고 thread.values, thread.toolCalls, thread.subgraphs, thread.output에 필요한 이벤트가 소비되지는 않습니다.

메시지 스트리밍

thread.messages를 채팅 모델 출력에 사용하세요.

async with client.threads.stream(assistant_id="agent") as thread:
    await thread.run.start(input=input)

    async for message in thread.messages:
        text = await message.text
        usage = (await message.output).usage_metadata
        print(text)
        print(usage)

message.text는 async 반복기이자 awaitable입니다. 반복하면 토큰 단위 출력, await하면 전체 텍스트가 나옵니다. message.reasoning은 추론 델타를, message.tool_calls(message.toolCalls)는 도구 호출 인자 청크를 노출합니다. message.output을 await하면 usage_metadata를 포함한 최종 메시지를 얻습니다. 텍스트·추론·도구 호출 청크를 도착 순서 그대로 소비해야 한다면 개별 투영 대신 원시 이벤트 스트림을 순회하세요.

상태 스트리밍

thread.values로 각 단계 후의 전체 상태 스냅샷을 스트리밍합니다.

async with client.threads.stream(assistant_id="agent") as thread:
    await thread.run.start(input=input)
    async for snapshot in thread.values:
        print(snapshot)
    final_state = await thread.output

thread.values 역시 awaitable입니다. await하면 최종 상태로 해석되며 await thread.output과 같습니다.

도구 호출 스트리밍

thread.tool_calls(thread.toolCalls)는 조립된 도구 호출을 노출합니다. 각 핸들은 도구 이름(call.name)과 조립된 입력(call.input, 일반 값—await하지 않음)을 지닙니다. 호출이 끝나면 call.output을 await해 도구 결과를 얻습니다.

async for call in thread.tool_calls:
    print(call.name, call.input)
    print(await call.output)
    if call.error is not None:
        print(call.error)

Python에서 call.deltas는 도구 출력이 스트리밍되는 대로 순회하는 async 반복기이고, 도구가 예외를 던졌다면 call.error에 예외가 담깁니다. 도구 이벤트는 도구 호출 ID로 thread.messages의 해당 도구 호출 콘텐츠 블록과 상관됩니다.

서브그래프 스트리밍

thread.subgraphs로 네임스페이스 문자열을 파싱하지 않고 중첩 그래프 작업을 관찰합니다.

async for subgraph in thread.subgraphs:
    print(subgraph.graph_name, subgraph.path)
    async for message in subgraph.messages:
        print(await message.text)

각 서브그래프 핸들은 그래프 이름(Python subgraph.graph_name / JS subgraph.name)과 네임스페이스 경로(Python subgraph.path / JS subgraph.namespace), 그리고 서브그래프별 messages·tool_calls·중첩 subgraphs 투영을 노출합니다. Deep Agents 배포에서는 서브에이전트 호출에 thread.subagents를 선호하세요—서브에이전트 이름과 서브에이전트별 메시지·도구 호출 투영을 제공합니다.

출력 스트리밍

run이 끝나면 thread.output을 await해 최종 상태를 얻습니다. thread.outputthread.values와 구독을 공유하므로, 다른 쪽도 읽고 있을 때 여분의 왕복이 필요하지 않습니다.

여러 투영 동시 스트리밍

애플리케이션 코드가 한 번에 두 개 이상의 투영을 필요로 하면 동시 소비자를 돌립니다.

import asyncio

async def consume_messages():
    async for message in thread.messages:
        print(await message.text)

async def consume_tool_calls():
    async for call in thread.tool_calls:
        print(call.name, await call.output)

async def consume_subgraphs():
    async for subgraph in thread.subgraphs:
        print(subgraph.graph_name, subgraph.path)

await asyncio.gather(consume_messages(), consume_tool_calls(), consume_subgraphs())

각 투영은 같은 스레드에 대한 필터링된 구독을 여므로, 동시 읽기는 실제 소비하는 채널 이상으로 서버 부하를 늘리지 않습니다.

인터럽트 후 재개

그래프가 인간 입력을 기다리며 일시정지하면 thread.interruptedthread.interrupts를 검사하고 인터럽트에 응답해 재개합니다.

async with client.threads.stream(assistant_id="agent") as thread:
    await thread.run.start(input=input)
    async for message in thread.messages:
        print(await message.text)

    if thread.interrupted:
        for interrupt in thread.interrupts:
            await thread.run.respond(
                {"decisions": [{"type": "approve"}]},
                interrupt_id=interrupt["interrupt_id"],
            )

    final_state = await thread.output

실행 중인 run 합류하기

페이지 리로드 후, 별도 워커에서, 또는 다른 클라이언트에서 이미 진행 중인 run에 붙으려면 기존 thread_id로 스레드 스트림을 열고 thread.run.start()를 건너뛰세요. 연결이 열릴 때 배포가 버퍼링된 이벤트를 재생(replay)하므로 소비자는 출력을 놓치지 않고 run 시작부터 상태를 재구성합니다.

모든 프로토콜 이벤트 스트리밍

애플리케이션 코드가 모든 이벤트를 필요로 하면 원시 프로토콜 이벤트 흐름을 읽습니다. Python은 thread.events를 순회하고, JavaScript는 subscribe를 엽니다(JS 스트림 객체 자체는 iterable이 아님).

async with client.threads.stream(assistant_id="agent") as thread:
    await thread.run.start(input=input)
    async for event in thread.events:
        print(event["method"], event["params"]["namespace"], event["params"]["data"])

특정 채널로 좁히려면 스레드에서 subscribe를 엽니다.

async for event in thread.subscribe(["messages", "tools"]):
    ...

각 이벤트는 채널별 페이로드를 감싼 ProtocolEvent 봉투입니다.

from typing import Any, NotRequired, TypedDict

class ProtocolEventParams(TypedDict):
    namespace: list[str]   # "<name>:<runtime_id>" 세그먼트 경로; []은 루트
    timestamp: int         # 밀리초 벽시계; 드리프트 가능하니 정렬에 의존하지 말 것
    data: Any              # 채널별 페이로드

class ProtocolEvent(TypedDict):
    type: str              # 항상 "event"
    seq: int               # 세션 내 증가; SSE id: 줄에 실림; 정렬에 사용
    method: str            # 채널 이름: "messages", "values", "tools", "lifecycle", "custom", ...
    params: ProtocolEventParams
    event_id: NotRequired[str]   # 세션 간 내구성 중복제거 ID, JSON 본문에 실림

namespace는 루트 그래프에서 이벤트를 발행한 스코프까지의 경로입니다. 루트는 빈 배열이고, 각 자식 실행이 "name:runtime_id" 세그먼트를 하나씩 추가하므로 서브그래프 안 중첩 도구 호출은 ["researcher:6f4d", "tools:91ac"]처럼 보입니다. 특정 하위 트리만 필요하면 원시 이벤트를 네임스페이스로 직접 필터링하세요. thread.subgraphs는 중첩 그래프 실행에 대해 이미 그렇게 합니다.

채널과 이벤트 라이프사이클

원시 이벤트는 채널 위로 흐릅니다. 채널 이름이 이벤트의 method로 나타나고, 각 채널은 고유한 이벤트 형태를 발행합니다.

채널 용도
values 전체 그래프 상태 스냅샷.
updates 노드별 상태 델타.
messages 콘텐츠 블록 중심의 채팅 모델 출력.
tools 도구 호출 시작, 스트리밍 출력, 종료, 오류 이벤트.
lifecycle run·서브그래프·서브에이전트 상태 변화.
checkpoints 브랜칭·타임트래블용 경량 체크포인트 봉투.
input 인간-루프 입력 요청과 응답.
tasks Pregel 태스크 생성·결과 이벤트.
custom 그래프 코드에서 온 사용자 정의 페이로드.
custom:<name> 애플리케이션 정의 스트림 트랜스포머 출력.

messages 채널은 출력을 콘텐츠 블록으로 모델링합니다. data.eventmessage-start, content-block-start, content-block-delta, content-block-finish, message-finish, error 중 하나입니다. 콘텐츠 블록은 명시적 경계를 갖습니다. 한 블록이 시작되고 0개 이상의 델타를 발행하며, 같은 메시지의 다음 블록이 시작되기 전에 끝납니다. message-finish는 토큰 사용량을 포함할 수 있고, 복구 불가능한 모델 호출 실패는 메시지 오류 이벤트로 도착합니다.

tools 채널은 도구 실행을 노출합니다. data.eventtool-started, tool-output-delta, tool-finished, tool-error 중 하나입니다. 도구 이벤트는 도구 호출 ID로 상관되므로 messages 채널의 원래 도구 호출 콘텐츠 블록에 다시 연결할 수 있습니다.

lifecycle 채널은 루트 run·서브그래프·서브에이전트 상태를 추적합니다. data.eventstarted, running, completed, failed, interrupted 중 하나입니다. 루트 run은 종료 상태 completed·failed·interrupted로 귀결되고, 자식 스코프는 running도 보고할 수 있습니다. 라이프사이클 데이터는 하위 스코프가 왜 시작됐는지(부모 도구 호출, fan-out send, 엣지 전환) 설명하는 선택적 graph_name·error·cause를 포함할 수 있습니다.

마지막 이벤트부터 재개

이벤트 스트림은 재개 가능합니다. Agent Server는 run별 이벤트를 유한 버퍼에 저장하고, 각각에 seq(세션 내 정렬)와 내구성 있는 event_id(재생·복제에 걸쳐 안정)를 부여하며, 재연결 시 커서부터 재생합니다. SDK는 일시적 드롭을 자동 처리합니다. 각 열린 구독이 관찰한 최고 seq를 추적하다가 재연결 시 그 커서부터 재생하고 event_id로 재전송된 이벤트를 중복제거합니다.

프로세스 경계 너머로 재개하려면—페이지 리로드, 워커 핸드오프, 별도 클라이언트—같은 thread_id로 스레드를 다시 엽니다. 새 구독이 열릴 때 서버가 버퍼링된 이벤트를 재생하고 SDK가 같은 타이핑된 투영으로 역다중화합니다. run별 버퍼가 유한하므로 매우 긴 run의 최초 이벤트는 쫓겨날 수 있습니다.

관련 자료

  • Streaming API — stream_mode 기반의 원시 스트리밍 API. langgraph-api>=0.10.0부터 지원.
  • LangGraph 이벤트 스트리밍 — 같은 개념을 인프로세스(in-process) LangGraph 애플리케이션에 적용.
  • LangSmith Deployment API — POST /threads/{thread_id}/stream/events 및 관련 엔드포인트 구현 수준 레퍼런스.

더 알아보기 (Learn more)

  • 실시간 출력의 또 다른 축인 일반 streaming(stream_mode)도 같은 허브의 관련 가이드를 참고하세요.
  • 이벤트를 소비하는 클라이언트 SDK 사용법은 LangGraph Python/JS SDK 문서에서 이어집니다.