StreamMux : le répartiteur d'événements pour flux multi-modes
Responsabilités
Dans le protocole de streaming v3 de LangGraph, lorsque l'utilisateur appelle graph.stream_events(version="v3"), il ne récupère plus une série de tuples bruts (namespace, mode, payload), mais un objet GraphRunStream qui expose des « projections » comme run.values / run.messages / run.updates / run.subgraphs / run.lifecycle (GraphRunStream:31-49). StreamMux est le centre de répartition qui transforme ces événements bruts en projections par mode, que l'utilisateur consomme ensuite via une simple boucle for. Le code se trouve dans libs/langgraph/langgraph/stream/_mux.py (StreamMux:26).
StreamMux n'implémente lui-même aucune logique de mode spécifique ; il ne fait que trois choses : enregistrer les transformers, pousser les événements, et se fermer. À chaque push, il appelle dans l'ordre d'enregistrement chaque StreamTransformer.process(event) (push:269-296), le transformer décide alors s'il injecte l'événement dans son canal (channel) de projection ; le mux (multiplexeur) ajoute ensuite l'événement au journal principal _events, qui constitue la « copie d'audit complète » du flux brut. Les projections ouvertes par les transformers sont toutes des instances de StreamChannel, et les canaux nommés sont automatiquement relayés par le mux vers le journal principal, de sorte que les effets de bord d'un transformer se manifestent à la fois sur son propre canal et dans le journal principal.
StreamMux est aussi la base de l'imbrication des sous-graphes. Quand un subgraph est découvert, SubgraphTransformer appelle mux._make_child(scope) (_make_child:193-225) pour créer un mini-mux pour le sous-graphe. Ce mini-mux partage le même ensemble de factories que le mux parent (donc une instance de transformer par scope), ainsi que la liaison pump du parent — dès que le pump racine avance, tous les événements des sous-graphes progressent d'un cran.
Motivation de conception
- pump piloté par l'appelant, pas de thread en arrière-plan : l'itérateur de
GraphRunStreamest le pump (_pump_next:107-130). Quand l'utilisateur faitfor event in run.values, chaque itération déclenche_pump_nextqui tire un événement brut pour le fournir au mux ; aucun thread n'est lancé en arrière-plan, la mémoire est dictée par la vitesse de consommation.bind_pumpenregistre cette fonction pump sur le mux (bind_pump:162-180), puis le mini-mux des sous-graphes l'hérite automatiquement via_make_child. - transformers en série dans l'ordre d'enregistrement :
pushutilisefor transformer in self._transformers: if not transformer.process(event): keep = False(push loop:288-292) — les effets de bord d'un transformer s'appliquent d'abord, le suivant les voit ; un transformer qui renvoieFalsesupprime seulement l'événement du journal principal, les transformers précédents l'ont déjà consommé. - distinction projection native vs extension : un transformer
_native = Truevoit sa projection attachée comme attribut direct (run.values) surGraphRunStream(native attrs:78-79) ; les non-natifs vont uniquement dansrun.extensions. L'ensemblenative_keys(native_keys:109) est la liste blanche « attributs directs » que le protocole v3 expose au SDK. before_builtinscède la place aux transformers de réécriture de contenu : les transformers de type filtrage PII ou modération de contenu doivent voir les champs texte bruts avant les transformers intégrés commeMessagesTransformer(before_builtins:94-109) ; à l'enregistrement, le mux réorganise en deux voies selonbefore_builtins = True(partition:131-148), mais l'ordre au sein de chaque voie est préservé.seqattribué seulement sur le mux racine :_assign_seqvautTruepar défaut, le mux enfant reçoitFalsevia_make_child(_assign_seq:215-220). Ainsi, les événements des sous-graphes transférés vers le journal racine ne voient pas leur enveloppe modifiée, etseqreste monotone par rapport à l'ordre d'écriture racine.
Fichiers clés
StreamMux 类:26-46— attributs et docstring du调度eur central.StreamMux.__init__:48-148— accepte une liste de transformers ou de factories, enregistre par voie selonbefore_builtins._make_child:193-225— crée un mini-mux à la découverte d'un sous-graphe, hérite de la liaison pump._register:227-267— appelletransformer.init()pour récupérer le dict de projection, vérifie les conflits de clés, lie les StreamChannel nommés.push:269-296— chemin synchrone : appelleprocessdans l'ordre, puis décide si l'événement entre dans le journal principal.apush:351-378— chemin asynchrone :aprocessen série await, garantit que les transformers suivants voient les résultats async des précédents.close:298-324— fermeture synchrone, appellefinalizesur chaque transformer, garantit la fermeture de tous les canaux même en cas d'erreur.aclose:380-404— fermeture asynchrone :gatherde toutes les tasks schedulées, puisafinalize._collect_stream_modes:398-414— déduit à partir desrequired_stream_modesdes transformers enregistrés les modes à demander au graphe._pregel_stream_v3:3533-3558— entrée synchrone v3, instancie StreamMux + GraphRunStream, repousse le stream_mode du mux versstream()sous-jacent.
Flux de données
Le cœur de StreamMux.push tient en quelques lignes — passer par les transformers dans l'ordre d'enregistrement, puis, si l'événement est toujours vivant, l'insérer dans le journal principal :
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) Ça paraît simple, mais transformer.process peut pousser dans son propre StreamChannel, et les canaux nommés sont automatiquement câblés par le mux vers le journal principal (_register:227-267) : un événement brut values peut donc être à la fois encapsulé par ValuesTransformer dans run.values et retransformé par un transformer utilisateur avant d'être réinjecté dans le journal principal. Les événements transitent en série entre les transformers, l'ordre est strictement préservé.
Limites et échecs
- transformer async en mode sync lève immédiatement :
_registerappelletransformer_requires_asyncpour vérifier (requires_async check:234-239) — si le transformer redéfinitaprocesssans mettresupports_sync=True, l'enregistrement sur un mux synchrone échoue et suggère d'utiliserastream()plutôt questream(). - conflit de clé de projection lève
ValueError:_registervérifieset(projection) & set(self.extensions)(conflict check:246-256) — le message d'erreur liste les clés en conflit et le transformer responsable ; un transformer personnalisé ne peut pas réutiliser un nom de clé d'un transformer intégré. _make_childrefuse les transformers pre-built : un mux construit avectransformers=ne peut pas créer de mini-mux (factory required:209-214), car les instances pré-construites ne peuvent pas être réinitialisées dans un nouveau scope.SubgraphTransformerne peut tourner que sur un mux construit avecfactories=.- erreur de finalize ne bloque pas le reste du nettoyage :
closeconserve le premier exception viafirst_error, continue à fermer les transformers et canaux restants, puis relâche à la fin (close error handling:312-324), garantissant qu'aucune ressource ne fuit. bind_pumpdoit être appelé avant la lecture du canal :wire_pump=Truepar défaut dansGraphRunStream.__init__déclenche_wire_request_more(_wire_request_more:80-92), le sous-mux hérite de la binding parent via_make_child; mais si l'utilisateur instancie lui-même un mux, oublierbind_pumplaisse_request_moreàNonedansStreamChannel.__iter__(_request_more:191-192) — l'itération se termine immédiatement.schedule(coro)réservé au mode async :StreamTransformer.schedulenécessite une event loop (schedule:233-263) — le muxgatherles tasks schedulées dansaclose, puisafinalizeaprès ; appelerscheduleen mode sync échoue.
Résumé
StreamMux est le pivot du protocole de streaming v3 : il n'implémente aucune logique de mode, il se contente d'ordonnancer les transformers, de créer des sous-mux par scope, et de collecter les événements bruts et les projections dans le journal principal. Pour le détail de chaque mode et la façon dont il transforme les événements en formes conviviales, voir stream/transformers ; pour la façon dont les événements bruts sont fournis par PregelLoop, voir moteur Pregel.
Voir la documentation officielle : LangGraph 文档 · README