Skip to content

Stream Transformers : transformer la sortie brute de Pregel en projections utilisables

源码版本1.2.9

Responsabilités

À chaque superpas (superstep), PregelLoop pousse via _emit un triplet (checkpoint_ns, mode, payload) dans une SyncQueue (_emit:1380-1414) — des événements bruts « valeurs », « mises à jour », « tâches », « points de contrôle (checkpoints) ». C'est la forme de sortie des protocoles v1 / v2 — l'utilisateur récupère une série de chunks. Mais en v3, l'utilisateur ne veut pas des chunks, il veut des « projections organisées par mode, qu'un seul consommateur peut itérer » : run.values est un flux de snapshots complets d'état, run.messages est un handle ChatModelStream sur le flux de tokens LLM, run.updates est un flux de mises à jour de nœuds, run.subgraphs est un handle vers l'entrée des sous-graphes.

La logique qui transforme les chunks bruts en ces projections se trouve dans libs/langgraph/langgraph/stream/transformers.py (ValuesTransformer:28) : un transformer par mode, tous héritent de la classe abstraite StreamTransformer (StreamTransformer:44). Ils sont enregistrés par StreamMux, alimentés en événements par GraphRunStream, et chacun maintient en interne un StreamChannel comme projection. Quand l'utilisateur appelle graph.stream_events(version="v3"), ce qu'il récupère réellement est un GraphRunStream (GraphRunStream:31), dont les attributs .values / .messages / .updates / .subgraphs / .custom / .lifecycle sont les points d'entrée des projections des transformers natifs (native attrs:78-79).

Le système de transformers est le point d'extension du protocole v3 — outre les cinq intégrés, l'utilisateur peut ajouter ses propres transformers via Pregel.compile(stream_transformers=[...]) (stream_transformers param:783), qui sont passés à StreamMux sous forme de factories dans _pregel_stream_v3 au côté des transformers intégrés (factories list:3533-3544).

Motivation de conception

  • Double voie projection native vs extension : un transformer _native = True voit sa projection attachée comme attribut direct run.<key> (_native:50-52) ; les non-natifs vont dans run.extensions, exposés seulement via le dictionnaire. Les cinq intégrés sont tous natifs ; les transformers utilisateurs vont par défaut dans extensions, pour obtenir un attribut direct il faut explicitement mettre _native = True.
  • required_stream_modes infère les modes à émettre : chaque transformer déclare les modes dont il a besoin (required_stream_modes:88-93), _collect_stream_modes prend l'union (_collect_stream_modes:398-414), et l'appel v3 à stream() sous-jacent passe cette union comme stream_mode (stream_mode=_collect_stream_modes:3549). L'utilisateur ne spécifie pas directement stream_mode, il « tire à la demande ».
  • Filtrage par scope pour éviter le bruit : ValuesTransformer.process vérifie params["namespace"] != self._scope_list (ValuesTransformer.process:70-82), n'accepte que les événements de son propre scope ; les événements des sous-graphes sont laissés au mini-mux du sous-graphe, le run.values racine n'est pas pollué par les graphes imbriqués.
  • Double voie sync / async + détection automatique : par défaut supports_sync = False ; si aprocess / afinalize / afail sont redéfinis, transformer_requires_async le classe comme async-only (transformer_requires_async:308-330), et l'enregistrement sur un mux synchrone échoue. Pour supporter sync et async, il faut explicitement supports_sync = True, comme SubgraphTransformer (supports_sync:689).
  • SubgraphTransformer utilise SubgraphRunStream plutôt que des événements bruts : à la découverte d'un appel de sous-graphe, _on_started crée un SubgraphRunStream qui enveloppe le mini-mux (SubgraphTransformer:670-705) — l'utilisateur récupère non pas un flux d'événements mais un handle objet, sur lequel il peut récursivement faire handle.values / handle.messages / handle.subgraphs, une interface isomorphe à celle du run racine.

Fichiers clés

Flux de données

Toute la chaîne v3 commence à PregelLoop._emit, passe par SyncQueue, GraphRunStream._pump_next, convert_to_protocol_event, et aboutit au mux. L'étape convert_to_protocol_event convertit le {type, ns, data, interrupts} v2 directement en ProtocolEvent v3 :

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) Notez que method est le type v2 — values / updates / messages / custom / tasks / checkpoints / debug, c'est la base du dispatch dans StreamTransformer.process, par exemple ValuesTransformer n'agit que sur les événements method == "values" (method check:70-71).

Le protocole v1 / v2 dans Pregel.stream suit une autre voie : les chunks bruts sortent directement de SyncQueue via la fonction _output (_output:4184-4243), qui selon que stream_mode est une chaîne ou une liste décide de yielder payload, (mode, payload) ou (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) C'est pourquoi stream_mode=["updates", "values"] donne des triplets alors qu'une chaîne unique donne une valeur unique — le même chunk brut est mis en forme à cet endroit dans _output. v3 ne passe pas par _output, mais par GraphRunStream._pump_nextmux.push → transformer.

Limites et échecs

  • process doit être implémentée, init doit retourner un dict : init est par défaut @abstractmethod, une sous-classe qui ne l'implémente pas échoue au démarrage (init abstract:127-138) ; process par défaut raise NotImplementedError (process default:162-164), ne pas la redéfinir fait exploser au premier push d'événement.
  • Ordre sensible entre transformers : les transformers before_builtins = True tournent d'abord (before_builtins:94-109) ; mais le docstring prévient qu'ils voient les événements tasks avant que LifecycleTransformer / SubgraphTransformer ne les consomment — modifier event["params"]["namespace"] ou les champs id / result / error / interrupts désordonne la comptabilité des deux derniers ; observez seulement, modifiez les champs froids.
  • MessagesTransformer couplé au protocole v2 : stream_events(version="v3") force version="v2" et CONFIG_KEY_STREAM_MESSAGES_V2=True lors de l'appel à stream() sous-jacent (_pregel_stream_v3 stream call:3545-3551), sinon la projection messages n'obtient pas le bon format de payload.
  • StreamChannel à consommateur unique, itérer deux fois lève une erreur : run.values ne peut être consommé que par une seule boucle for (single-consumer note:38-40) ; pour fan-out il faut projection.tee(n) (tee:245), au prix d'un tampon.
  • schedule(on_error="raise") transforme le chemin de fermeture en échec : on_error="raise" d'un task programmé pousse le mux vers afail dans aclose (on_error:254-257) ; "log" par défaut se contente de journaliser sans propager — un mauvais choix fait qu'un seul transformer qui échoue emporte toute la chaîne.
  • SubgraphTransformer._should_track strictement supérieur au scope : ne traque que les sous-graphes au niveau immédiatement inférieur (should_track:634-636) ; les sous-sous-graphes sont découverts récursivement via handle.subgraphs du handle enfant ; si l'utilisateur ignore handle.subgraphs et regarde directement run.subgraphs racine, les événements petits-enfants n'y apparaîtront pas.

Résumé

StreamTransformer est le point d'extension du protocole de streaming v3 : cinq transformers intégrés transforment les (ns, mode, payload) bruts en projections typées itérables, et l'utilisateur peut ajouter ses propres transformers via le paramètre stream_transformers. Pour tout le mécanisme de répartition, voir StreamMux ; pour la façon dont les chunks bruts sont alimentés depuis PregelLoop dans SyncQueue, voir PregelLoop.

Voir la documentation officielle : LangGraph 文档 · README