Skip to content

StreamMux : le répartiteur d'événements pour flux multi-modes

源码版本1.2.9

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 GraphRunStream est le pump (_pump_next:107-130). Quand l'utilisateur fait for event in run.values, chaque itération déclenche _pump_next qui 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_pump enregistre 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 : push utilise for 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 renvoie False supprime 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 = True voit sa projection attachée comme attribut direct (run.values) sur GraphRunStream (native attrs:78-79) ; les non-natifs vont uniquement dans run.extensions. L'ensemble native_keys (native_keys:109) est la liste blanche « attributs directs » que le protocole v3 expose au SDK.
  • before_builtins cè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 comme MessagesTransformer (before_builtins:94-109) ; à l'enregistrement, le mux réorganise en deux voies selon before_builtins = True (partition:131-148), mais l'ordre au sein de chaque voie est préservé.
  • seq attribué seulement sur le mux racine : _assign_seq vaut True par défaut, le mux enfant reçoit False via _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, et seq reste 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 selon before_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 — appelle transformer.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 : appelle process dans l'ordre, puis décide si l'événement entre dans le journal principal.
  • apush:351-378 — chemin asynchrone : aprocess en série await, garantit que les transformers suivants voient les résultats async des précédents.
  • close:298-324 — fermeture synchrone, appelle finalize sur chaque transformer, garantit la fermeture de tous les canaux même en cas d'erreur.
  • aclose:380-404 — fermeture asynchrone : gather de toutes les tasks schedulées, puis afinalize.
  • _collect_stream_modes:398-414 — déduit à partir des required_stream_modes des 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 vers stream() 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 :

python
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 : _register appelle transformer_requires_async pour vérifier (requires_async check:234-239) — si le transformer redéfinit aprocess sans mettre supports_sync=True, l'enregistrement sur un mux synchrone échoue et suggère d'utiliser astream() plutôt que stream().
  • conflit de clé de projection lève ValueError : _register vérifie set(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_child refuse les transformers pre-built : un mux construit avec transformers= 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. SubgraphTransformer ne peut tourner que sur un mux construit avec factories=.
  • erreur de finalize ne bloque pas le reste du nettoyage : close conserve le premier exception via first_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_pump doit être appelé avant la lecture du canal : wire_pump=True par défaut dans GraphRunStream.__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, oublier bind_pump laisse _request_more à None dans StreamChannel.__iter__ (_request_more:191-192) — l'itération se termine immédiatement.
  • schedule(coro) réservé au mode async : StreamTransformer.schedule nécessite une event loop (schedule:233-263) — le mux gather les tasks schedulées dans aclose, puis afinalize après ; appeler schedule en 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