Stream Transformers:生の Pregel 出力を利用可能な投影にする
役割
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.processはparams["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_startedはSubgraphRunStreamを建てて 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-330—aprocess/afinalize/afailがオーバーライドされているか検査し、transformer が sync mux で走れるか決めます。ValuesTransformer:28-82—run.values投影。required_stream_modes=("values",)。UpdatesTransformer:120-152—run.updates投影。MessagesTransformer:155-335—run.messages投影。LLM トークンストリームをChatModelStreamハンドルに詰め、v2 プロトコルの(payload, metadata)入力を処理します。LifecycleTransformer:608-667—run.lifecycle投影。サブグラフライフサイクルイベントを native protocol event として主ログに戻し、リモート SDK から見えるようにします。SubgraphTransformer:670-705—run.subgraphs投影。サブグラフ呼び出しを発見したときSubgraphRunStream+ mini-mux を建てます。GraphRunStream._pump_next:107-130— ポンプ:chunk を 1 つ引き、convert_to_protocol_eventを通して mux に投げます。convert_to_protocol_event:10-32— v2StreamPart→ v3ProtocolEventのフィールドマッピング。
データフロー
v3 のリンク全体は PregelLoop._emit から始まり、SyncQueue、GraphRunStream._pump_next、convert_to_protocol_event を経て mux に至ります。convert_to_protocol_event は v2 の {type, ns, data, interrupts} を v3 ProtocolEvent に直接変換します:
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 でのディスパッチ根拠です。例えば ValuesTransformer は method == "values" のイベントだけを相手にします (method check:70-71)。
一方 v1 / v2 プロトコルは Pregel.stream で別の経路を走ります。生 chunk は SyncQueue から直接出て _output 関数を通ります (_output:4184-4243)。stream_mode が文字列か list かで payload か (mode, payload) か (ns, mode, payload) のいずれかを出力します:
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_next → mux.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 に変える:scheduleのon_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