커스텀 스트림 채널

커스텀 스트림 채널 (Custom stream channels)

LangGraph 에이전트는 메시지와 툴 호출 그 이상을 스트리밍할 수 있어요. 서버 쪽 **스트림 트랜스포머(stream transformer)**가 클라이언트로 흐르는 프로토콜을 검사·재작성하고, 이름이 붙은 **커스텀 채널(custom channel)**에 자기만의 구조화 데이터를 게시할 수 있죠. 이 글에서는 고객지원 에이전트가 브라우저에 도달하기 전에 모든 이벤트에서 PII(이메일, 전화번호, SSN, 카드 번호, IP)를 지우고, 수정 횟수를 redaction-stats 채널로 게시하는 예시를 중심으로 다룰게요.

출처: 공식문서

커스텀 채널이 동작하는 방식

커스텀 채널은 두 끝단이 있어요. 서버에서는 StreamTransformer가 이름이 붙은 StreamChannel을 열고 여기에 페이로드를 밀어 넣어요. 클라이언트에서는 셀렉터가 그에 맞는 custom:<name> 채널을 구독해 페이로드를 반응형 상태로 노출하죠.

트랜스포머의 process 메서드는 모든 프로토콜 이벤트마다 실행돼요. 이벤트를 그 자리에서 변형할 수 있고(여기서는 messages, tools, values 데이터에서 PII를 지움), 보고할 게 생길 때마다 사이드 채널 업데이트를 푸시해요. 클라이언트 셀렉터(useExtension, useChannel)는 v1 프론트엔드 SDK 패키지(@langchain/react, @langchain/vue, @langchain/svelte, @langchain/angular)에 포함돼 있어요.

import time

from langgraph.stream import ProtocolEvent, StreamChannel, StreamTransformer


class RedactionStatsTransformer(StreamTransformer):
    def __init__(self, scope: tuple[str, ...] = ()) -> None:
        super().__init__(scope)
        # Open a channel named "redaction-stats".
        self.redaction_stats = StreamChannel("redaction-stats")
        self.counts = empty_counts()

    def init(self) -> dict[str, StreamChannel]:
        return {"redactionStats": self.redaction_stats}

    def process(self, event: ProtocolEvent) -> bool:
        # Redact event["params"]["data"] in place and tally what was found.
        delta = redact_in_place(event, self.counts)
        if delta:
            # Publish a payload on the channel.
            self.redaction_stats.push(
                {
                    "kind": "update",
                    "at": int(time.time() * 1000),
                    "delta": delta,
                    "counts": dict(self.counts),
                    "total": sum(self.counts.values()),
                }
            )
        return True  # Keep the (now-redacted) event in the stream.


def create_redaction_stats_transformer() -> RedactionStatsTransformer:
    return RedactionStatsTransformer()

에이전트를 만들 때 트랜스포머를 붙여요:

from langchain.agents import create_agent

agent = create_agent(
    model="anthropic:claude-haiku-4-5",
    tools=[...],
    transformers=[create_redaction_stats_transformer],
)

useStream 설정하기

useStream은 평소처럼 연결하면 돼요. 커스텀 채널 셀렉터는 여기서 반환된 같은 stream 핸들을 받아요.

최신 페이로드 읽기 — useExtension

useExtensioncustom:<name> 채널을 구독해 트랜스포머가 가장 최근에 푸시한 페이로드를 이미 풀어서(unwrapped) 타입이 붙은 채로 돌려줘요. UI가 현재 값만 필요할 때(라이브 카운터, 진행률, 상태 배지) 가장 편리한 선택이에요. 채널 이름("redaction-stats")만 넘기면 되고 custom: 접두사는 붙이지 않아요.

반환값은 각 프레임워크의 반응형 모델을 따르는데, React·Svelte에서는 일반 값, Vue에서는 Ref(latest.value), Angular에서는 시그널(latest())이에요. 첫 페이로드가 도착하기 전까지는 값이 undefined예요. 선택적인 세 번째 target 인자는 useMessages(stream, node)가 발견된 그래프 노드로 메시지를 범위 지정하는 것과 같은 방식으로, 네임스페이스로 구독을 범위 지정해요. 네임스페이스 타겟팅은 Graph execution을 참고하세요.

원시 이벤트 버퍼링 — useChannel

useChannel은 원시 이벤트 탈출구예요. 하나 이상의 채널을 구독하고 단일 언랩 값 대신 원시 프로토콜 이벤트의 제한된 버퍼(bounded buffer)를 돌려줘요. 최신 값이 아니라 **이력(history)**이 필요할 때(이벤트 로그, 감사 추적)나 상위 레벨 셀렉터가 다루지 않는 채널이 필요할 때 사용해요. 전체 채널 id("custom:redaction-stats")를 전달해요.

각 항목은 원시 프로토콜 이벤트이므로 페이로드는 event.params.data 아래에 있어요. 직접 풀어야 해요:

function parseRedactionStatsEvents(rawEvents: Event[]): RedactionStatsEvent[] {
  const out: RedactionStatsEvent[] = [];
  for (const event of rawEvents) {
    const data = event.params?.data;
    const payload = data?.payload ?? data;
    if (payload?.kind === "update") out.push(payload);
  }
  return out;
}

옵션 인자로 버퍼를 제어해요:

const rawEvents = useChannel(
  stream,
  ["custom:redaction-stats"],
  undefined, // target namespace
  { bufferSize: 200, replay: true },
);
Option Default 효과
bufferSize "default" 버퍼링하는 이벤트 최대 개수. 한도를 넘으면 오래된 이벤트부터 버려져요.
replay true 셀렉터가 마운트될 때 채널에서 이미 본 이벤트를 재생해요. 라이브 이벤트만 보려면 끄면 돼요.

useExtension vs useChannel 선택

둘 다 같은 커스텀 채널을 읽지만 반환하는 내용은 달라요.

useExtension useChannel
반환값 최신 페이로드 (T | undefined) 원시 이벤트의 제한된 버퍼 (Event[])
형태 언랩되고 타입이 붙은 페이로드 원시 프로토콜 이벤트; event.params.data를 직접 언랩
구독 기준 채널 이름 ("redaction-stats") 전체 채널 id (["custom:redaction-stats"])
사용 시점 현재 값이 필요할 때 이력·로그·여러 채널이 필요할 때
옵션 bufferSize, replay

흔한 패턴은 같은 채널에 둘을 함께 쓰는 거예요. useExtension으로 라이브 요약(현재 합계)을, useChannel로 스레드 전체에 걸친 매 업데이트의 스크롤 이벤트 로그를 지원하죠.

사용 사례

커스텀 채널은 메시지·툴 호출·그래프 상태로 깔끔하게 매핑되지 않는 서버 쪽 신호에 어울려요.

  • 컴플라이언스·레드액션 통계: 지워진 PII, 차단된 콘텐츠, 정책 적중 횟수(위 예시처럼).
  • 진행 보고: 오래 걸리는 툴이 내는 진행률(%)이나 단계 라벨.
  • 라이브 메트릭: 실행 중 누적되는 토큰 사용량, 지연 시간, 비용.
  • 소스·인용: 에이전트가 답을 뒷받침할 때 검색한 문서를 사이드 패널로 푸시.
  • 도메인 이벤트: 메시지 기록을 바꾸지 않고 백엔드가 표면화하고 싶은 구조화 업데이트.

더 알아보기 (Learn more)

  • Overview — LangGraph 프론트엔드 스트림 API와 아키텍처
  • Graph execution — 멀티 노드 파이프라인을 위한 네임스페이스 범위 셀렉터