StreamMux: Event-Verteiler für Multi-Modus-Streams
Verantwortung
Im v3-Streaming-Protokoll von LangGraph erhält der Nutzer bei graph.stream_events(version="v3") keine Rohfolge von (namespace, mode, payload)-Tupeln mehr, sondern ein GraphRunStream-Objekt mit Projektionen (projection) wie run.values / run.messages / run.updates / run.subgraphs / run.lifecycle (GraphRunStream:31-49). Diese Roh-Events nach Modus auf verschiedene Projektionen aufzuteilen und dem Nutzer über eine gewöhnliche for-Schleife zugänglich zu machen, ist Aufgabe des StreamMux – der Event-Verteilungszentrale. Der Code liegt in libs/langgraph/langgraph/stream/_mux.py (StreamMux:26).
StreamMux selbst implementiert keine modusspezifische Logik; es macht nur drei Dinge: Transformer registrieren, Events push-en und schließen. Beim Push wird für jeden registrierten StreamTransformer.process(event) in Registrierreihenfolge aufgerufen (push:269-296); der Transformer entscheidet, ob das Event in seinen eigenen projection-Channel aufgenommen wird. Danach hängt der Mux das Event an _events an – das Hauptlog ist eine «vollständige Audit-Kopie» des Roh-Streams. Die Projektionen, die ein Transformer öffnet, sind StreamChannel-Instanzen; benannte Channels werden vom Mux zusätzlich automatisch ins Hauptlog weitergeleitet, sodass die Seiteneffekte eines Transformers sowohl in seinem eigenen Channel als auch im Hauptlog sichtbar sind.
StreamMux ist außerdem die Grundlage für die Schachtelung von Untergraphen (subgraph). Sobald ein Subgraph entdeckt wird, ruft SubgraphTransformer mux._make_child(scope) auf (_make_child:193-225), um einen Mini-Mux für den Subgraphen zu erzeugen. Dieser Mini-Mux teilt sich denselben Satz an Factorys mit dem Eltern-Mux (jeder Transformer ist also pro Scope einmal vorhanden) und auch die Pump-Bindung des Eltern-Mux – sobald die Root-Pump sich bewegt, rücken alle Subgraph-Events um eine Position nach vorne.
Entwurfsmotivation
- Aufrufergetriebene Pump, kein Hintergrund-Thread: Der Iterator von
GraphRunStreamist die Pump (_pump_next:107-130). Wenn der Nutzerfor event in run.valuesausführt, löst jede Iteration_pump_nextaus und holt ein Roh-Event, das dem Mux zugeführt wird; es wird kein Hintergrund-Thread gestartet, und der Speicherverbrauch richtet sich nach der Konsumrate.bind_pumpregistriert diese Pump-Funktion im Mux (bind_pump:162-180) und wird über_make_childautomatisch vererbt. - Transformer in Registrierreihenfolge seriell:
pushnutztfor transformer in self._transformers: if not transformer.process(event): keep = False(push loop:288-292); die Seiteneffekte des vorigen Transformers finden zuerst statt, und der nächste sieht sie bereits; ein Transformer, derFalsezurückgibt, unterdrückt das Event nur im Hauptlog, die vorherigen Transformer haben es bereits konsumiert. - native vs. extension Projektionen: Transformer mit
_native = Truewerden vom Mux als direkte Attribute wierun.valuesanGraphRunStreamangehängt (native attrs:78-79); nicht-native landen nur inrun.extensions. Die Mengenative_keys(native_keys:109) ist die Whitelist an «direkten Attributen», die das v3-Protokoll dem SDK freigibt. before_builtinsmacht Platz für inhaltsmodifizierende Transformer: Transformer für PII-Filterung oder Inhaltsmoderation müssen die rohen Textfelder vor den eingebauten Transformern wieMessagesTransformersehen (before_builtins:94-109); der Mux teilt die Registrierung in zwei Lanes auf, wennbefore_builtins = True(partition:131-148), aber die Reihenfolge innerhalb einer Lane bleibt erhalten.- seq wird nur im Root-Mux vergeben:
_assign_seqist per Default True, der Child-Mux erhält über_make_childden WertFalse(_assign_seq:215-220). So werden Subgraph-Events beim direkten Weiterleiten ins Root-Log nicht mit neuem Umschlag versehen, undseqbleibt monoton und entspricht der Schreibreihenfolge im Root.
Schlüsseldateien
StreamMux Klasse:26-46— Zentrale Scheduler-Attribute und Docstring.StreamMux.__init__:48-148— Nimmt eine Transformer-Liste oder Factories, registriert nachbefore_builtinsin zwei Lanes._make_child:193-225— Erzeugt beim Entdecken eines Subgraphen einen Mini-Mux und vererbt die Pump-Bindung._register:227-267— Rufttransformer.init()auf, um das projection-Dict zu holen, prüft Schlüsselkonflikte und bindet benannte StreamChannels.push:269-296— Synchroner Pfad: ruftprocessin Reihenfolge auf und entscheidet dann über die Aufnahme ins Hauptlog.apush:351-378— Asynchroner Pfad:aprocesswird seriell awaited, sodass nachfolgende Transformer die async-Ergebnisse der vorherigen sehen.close:298-324— Synchrones Schließen; ruftfinalizejedes Transformers auf; im Fehlerfall werden dennoch alle Channels geschlossen.aclose:380-404— Asynchrones Schließen; zunächstgatheraller geschedulten Tasks, dannafinalize._collect_stream_modes:398-414— Leitet aus denrequired_stream_modesder registrierten Transformer ab, welche Modi beim Graphen angefordert werden müssen._pregel_stream_v3:3533-3558— V3-Synchroner Einstieg; instanziiert StreamMux + GraphRunStream und reicht den vom Mux bestimmte stream_mode an das zugrunde liegendestream()weiter.
Datenfluss
Der Kern von StreamMux.push besteht nur aus wenigen Zeilen – die Transformer werden in Registrierreihenfolge durchlaufen, und wenn das Event überlebt, wird es ins Hauptlog eingefügt:
def push(self, event: ProtocolEvent) -> None:
keep = True
for transformer in self._transformers:
if not transformer.process(event):
keep = False
if keep:
if self._assign_seq:
self._seq += 1
event["seq"] = self._seq
self._events.push(event)(push body:288-296) Es wirkt unscheinbar, aber in transformer.process kann ein Event in den eigenen StreamChannel gepusht werden, und benannte Channels werden vom Mux automatisch ins Hauptlog gewired (_register:227-267). So kann ein Roh-values-Event gleichzeitig von ValuesTransformer in run.values übernommen und von einem nutzerdefinierten Transformer (transformer) transformiert und zurück ins Hauptlog geschrieben werden. Events laufen über eine serielle Kette zwischen Transformern, die Reihenfolge bleibt strikt erhalten.
Grenzen und Fehler
- Async Transformer im sync-Modus wirft sofort:
_registerrufttransformer_requires_asyncauf (requires_async check:234-239); wenn ein Transformeraprocessüberschreibt, aber nichtsupports_sync=Truesetzt, schlägt die Registrierung am sync-Mux fehl und fordertastream()stattstream(). - projection-Schlüsselkonflikt wirft direkt
ValueError:_registerprüftset(projection) & set(self.extensions)(conflict check:246-256) und listet im Fehler die Konfliktschlüssel sowie die besitzenden Transformer auf; ein eigener Transformer darf keine Projection-Keys verwenden, die mit eingebauten kollidieren. _make_childlehnt vorgefertigte Transformer ab: Wird der Mux mit dem Parametertransformers=konstruiert, kann kein Child-Mux erzeugt werden (factory required:209-214), da vorgefertigte Instanzen sich in einem neuen Scope nicht neu initialisieren lassen.SubgraphTransformerläuft nur auf einem Mux, der mitfactories=konstruiert wurde.- Fehler in finalize blockiert nicht die übrige Aufräumung:
closefängt die erste Ausnahme infirst_errorund schließt trotzdem alle weiteren Transformer und Channels, bevor die Ausnahme am Ende erneut geworfen wird (close error handling:312-324), damit keine Ressourcen lecken. bind_pumpmuss vor dem Lesen des Channels aufgerufen werden:GraphRunStream.__init__mitwire_pump=Trueruft per Default_wire_request_moreauf (_wire_request_more:80-92); der Child-Mux übernimmt beim_make_childautomatisch die Eltern-Bindung. Wer einen Mux selbst instanziiert und vergisst, die Pump zu binden, lässtStreamChannel.__iter__'s_request_morefür immerNonesein (_request_more:191-192), und die Iteration endet sofort.schedule(coro)funktioniert nur im async-Modus:StreamTransformer.schedulebenötigt einen Event-Loop (schedule:233-263); der Mux sammelt inaclosealle geschedulten Tasks mitgatherein und ruft erst nach Abschlussafinalizeauf; im sync-Modus schlägt schedule fehl.
Zusammenfassung
StreamMux ist der Knotenpunkt des v3-Streaming-Protokolls: Es implementiert selbst keine Modus-Logik, sondern ruft Transformer in Reihenfolge auf, baut pro Scope einen Child-Mux auf und nimmt sowohl Roh-Events als auch Transformer-Projektionen ins Hauptlog auf. Wie jeder Modus die Events in eine nutzerfreundliche Form wandelt, siehe stream/transformers; wie die Roh-Events aus dem PregelLoop eingespeist werden, siehe Pregel 引擎.
Siehe offizielle Dokumentation: LangGraph 文档 · README.