Skip to content

Stream Transformers:生の Pregel 出力を利用可能な投影にする

源码版本1.2.9

役割

PregelLoop は各スーパーステップで _emit を通じて (checkpoint_ns, mode, payload) の 3-tuple を SyncQueue に push します (_emit:1380-1414)。中身は生の「値」「更新」「タスク」「チェックポイント」といったイベントです。これは v1 / v2 プロトコルの出力形態で、ユーザーが受け取るのは chunk の羅列です。しかし v3 プロトコルでは、ユーザーは chunk ではなく「モード別に組織され、単一消費者でイテレート可能な投影 (projection)」を希望します:run.values は完全状態スナップショットのストリーム、run.messages は LLM トークンストリームの ChatModelStream ハンドル、run.updates はノード更新ストリーム、run.subgraphs はサブグラフ入口ハンドルです。

生 chunk をこれらの投影に変える論理が libs/langgraph/langgraph/stream/transformers.py にあります (ValuesTransformer:28)。モードごとに 1 つの transformer で、共通に StreamTransformer 抽象基底クラス (StreamTransformer:44) を継承します。これらは StreamMux に登録され、GraphRunStream からイベントを push され、内部に StreamChannel を投影として保持します。ユーザーが graph.stream_events(version="v3") を呼ぶと、実際には GraphRunStream (GraphRunStream:31) を受け取り、その .values / .messages / .updates / .subgraphs / .custom / .lifecycle の各プロパティが native transformer の投影入口です (native attrs:78-79)。

transformer 体系全体が v3 プロトコルの拡張点です。組込みの 5 つに加え、ユーザーは Pregel.compile(stream_transformers=[...]) で独自 transformer を追加でき (stream_transformers param:783)、_pregel_stream_v3 で組込み transformer と一緒に factory 形式で StreamMux に渡されます (factories list:3533-3544)。

設計動機

  • native vs extension 投影の二系統:_native = True の transformer 投影は run.<key> の直接属性にぶら下がり (_native:50-52)、非 native は run.extensions に入り辞書としてだけ晒されます。組込みの 5 つはすべて native で、ユーザーカスタムはデフォルトで extensions に入ります。直接属性にするには明示的に _native = True を設定します。
  • required_stream_modes でグラフが発するモードを逆推:transformer 自身が必要な mode を宣言し (required_stream_modes:88-93)、_collect_stream_modes が和集合を取り (_collect_stream_modes:398-414)、v3 が下層の stream() を呼ぶときにこの和集合を stream_mode として渡します (stream_mode=_collect_stream_modes:3549)。ユーザーは stream_mode を直接指定せず、「必要に応じて引く」形です。
  • scope フィルタでノイズを回避:ValuesTransformer.processparams["namespace"] != self._scope_list を検査し (ValuesTransformer.process:70-82)、自身の scope のイベントだけを受け付けます。サブグラフのイベントはサブグラフ自身の mini-mux に処理を委ね、ルートの run.values がネストされたグラフに汚染されないようにします。
  • sync / async 二系統 + 自動検出:transformer はデフォルト supports_sync = False で、aprocess / afinalize / afail をオーバーライドすると transformer_requires_async に async-only と判定され (transformer_requires_async:308-330)、sync mux への登録は即座にエラーになります。sync / async 両対応するには明示的に supports_sync = True にする必要があります。例えば SubgraphTransformer (supports_sync:689)。
  • SubgraphTransformer は裸イベントではなく SubgraphRunStream を渡す:サブグラフ呼び出しを発見したとき、_on_startedSubgraphRunStream を建てて mini-mux を包み (SubgraphTransformer:670-705)、ユーザーが受け取るのはイベントストリームではなくハンドルオブジェクトです。ハンドル上で handle.values / handle.messages / handle.subgraphs を再帰的にたどれ、ルート run と同型のインタフェースです。

主要ファイル

  • StreamTransformer 基类:44-115 — 抽象 init / process / aprocess / finalize / fail + ClassVar 設定。
  • init:128-138 — projection dict を返す。key は extensions に入り、_native=True のときは run.<key> にもぶら下がります。
  • process:148-164 — sync レーン。デフォルトは raise NotImplementedError で、サブクラスがオーバーライド必須です。
  • aprocess:166-184 — async レーン。デフォルトは process に委譲。async 仕事が必要ならオーバーライドします。
  • schedule:233-263 — mux ライフサイクルに紐づく asyncio.Task。aclose 時に gather し、afail 時に cancel します。
  • transformer_requires_async:308-330aprocess / afinalize / afail がオーバーライドされているか検査し、transformer が sync mux で走れるか決めます。
  • ValuesTransformer:28-82run.values 投影。required_stream_modes=("values",)
  • UpdatesTransformer:120-152run.updates 投影。
  • MessagesTransformer:155-335run.messages 投影。LLM トークンストリームを ChatModelStream ハンドルに詰め、v2 プロトコルの (payload, metadata) 入力を処理します。
  • LifecycleTransformer:608-667run.lifecycle 投影。サブグラフライフサイクルイベントを native protocol event として主ログに戻し、リモート SDK から見えるようにします。
  • SubgraphTransformer:670-705run.subgraphs 投影。サブグラフ呼び出しを発見したとき SubgraphRunStream + mini-mux を建てます。
  • GraphRunStream._pump_next:107-130 — ポンプ:chunk を 1 つ引き、convert_to_protocol_event を通して mux に投げます。
  • convert_to_protocol_event:10-32 — v2 StreamPart → v3 ProtocolEvent のフィールドマッピング。

データフロー

v3 のリンク全体は PregelLoop._emit から始まり、SyncQueue、GraphRunStream._pump_next、convert_to_protocol_event を経て mux に至ります。convert_to_protocol_event は v2 の {type, ns, data, interrupts} を v3 ProtocolEvent に直接変換します:

python
def convert_to_protocol_event(part: StreamPart) -> ProtocolEvent:
    part_dict = cast(dict[str, Any], part)
    params: _ProtocolEventParams = {
        "namespace": list(part_dict["ns"]),
        "timestamp": int(time.time() * 1000),
        "data": part_dict["data"],
    }
    if "interrupts" in part_dict:
        params["interrupts"] = part_dict["interrupts"]
    return {
        "type": "event",
        "method": part_dict["type"],
        "params": params,
    }

(convert body:20-32) method は v2 の type です。values / updates / messages / custom / tasks / checkpoints / debug で、これが StreamTransformer.process でのディスパッチ根拠です。例えば ValuesTransformermethod == "values" のイベントだけを相手にします (method check:70-71)。

一方 v1 / v2 プロトコルは Pregel.stream で別の経路を走ります。生 chunk は SyncQueue から直接出て _output 関数を通ります (_output:4184-4243)。stream_mode が文字列か list かで payload(mode, payload)(ns, mode, payload) のいずれかを出力します:

python
if version == "v2":
    ...
    yield {"type": mode, "ns": ns, "data": payload, "interrupts": ints}
elif stream_subgraphs and isinstance(stream_mode, list):
    yield (ns, mode, payload)
elif isinstance(stream_mode, list):
    yield (mode, payload)
elif stream_subgraphs:
    yield (ns, payload)
else:
    yield payload

(_output shape branches:4220-4243) だから stream_mode=["updates", "values"] は 3-tuple を取得し、単一文字列は単値を取得します。同じ生 chunk が _output で形状を分けられます。v3 は _output を通らず、GraphRunStream._pump_nextmux.push → transformer を走ります。

境界と失敗

  • process は 1 つ実装必須、init は dict を返す必要あり:init はデフォルトで @abstractmethod で、サブクラスが実装しないと起動時に失敗します (init abstract:127-138)。process はデフォルト raise NotImplementedError (process default:162-164) で、オーバーライドしないと最初のイベント push 時に炸裂します。
  • transformer 間で順序センシティブ:before_builtins = True の transformer が先に走り (before_builtins:94-109) ますが、docstring は tasks イベントが LifecycleTransformer / SubgraphTransformer に消費される前に見えると警告します。event["params"]["namespace"]id / result / error / interrupts フィールドを変更すると後二者の記帳が乱れるため、観察だけにし、冷たいフィールドだけを変えるべきです。
  • MessagesTransformer は v2 プロトコルと結合:stream_events(version="v3") は下層の stream() を呼ぶとき強制的に version="v2"CONFIG_KEY_STREAM_MESSAGES_V2=True を渡します (_pregel_stream_v3 stream call:3545-3551)。さもなくば messages 投影が正しい payload フォーマットを取得できません。
  • StreamChannel は単一消費者で、2 回イテレートするとエラー:run.values は 1 つの for ループでしか消費できません (single-consumer note:38-40)。fan-out するには projection.tee(n) (tee:245) が必要で、代償はバッファです。
  • schedule(on_error="raise") は close パスを fail に変える:scheduleon_error="raise" は task 异常が aclose 時に mux を afail に押しやります (on_error:254-257)。デフォルトの "log" はログだけ取り伝播しません。選び方を間違えると単一 transformer の失敗がストリーム全体を引きずり下ろします。
  • SubgraphTransformer._should_track は厳格に scope より深い層だけ:現在の scope の 1 層下のサブグラフだけを追跡し (should_track:634-636)、孫グラフは子ハンドルの handle.subgraphs で再帰的に発見します。ユーザーが handle.subgraphs を飛ばしてルートの run.subgraphs を直接見た場合、孫グラフのイベントはそこには現れません。

まとめ

StreamTransformer は v3 ストリーミングプロトコルの拡張点です。5 つの組込み transformer が生の (ns, mode, payload) をイテレート可能な型付き投影に変え、ユーザーも自身の transformer を stream_transformers パラメータに追加できます。ディスパッチ機構全体は StreamMux を、生 chunk が PregelLoop からどう SyncQueue に流れ込むかは PregelLoop を参照してください。

公式資料:LangGraph 文档 · README