Skip to content

Streaming de sortie : astream et stream_mode

源码版本1.2.9

Responsabilités

Pregel.astream / Pregel.stream sont les points d'entrée d'exécution exposés par le graphe compilé. Vous y alimentez vos entrées, ils démarrent un PregelLoop qui tourne tours après tours, et en cours de route projettent les données produites à chaque superpas (superstep) sous forme d'événements de flux selon le stream_mode spécifié par l'appelant, events vers une file d'attente, puis l'appelant les consomme par itération. ainvoke / invoke sont essentiellement astream / stream convergés en une valeur unique — le corps de Pregel.invoke est for chunk in self.stream(...) en prenant le dernier fragment (Pregel.invoke body:3891).

La responsabilité de cette couche est de découpler « exécution » et « forme de sortie » : l'exécution reste le cycle BSP de Pregel, la forme de sortie est déterminée par stream_mode — en un seul run du graphe, vous pouvez demander values (état complet à chaque étape), updates (incrémentiel par étape), messages (flux de tokens LLM), custom (sortie personnalisée intra-nœud), checkpoints (événements de point de contrôle / checkpoint), tasks (début/fin de tâche), debug (tout), voire passer une liste pour récupérer plusieurs modes simultanément.

Motivation de conception

  • Ne pas bloquer l'exécution : la sortie passe par une file, runner.tick / atick y pousse un événement dès qu'une tâche est terminée, la boucle for externe l'iterate à mesure — au lieu d'attendre la fin complète du graphe pour tout renvoyer d'un coup. C'est indispensable pour le flux de tokens LLM (messages) — l'utilisateur doit voir la génération en direct.
  • Multi-projection réutilise la même exécution : au sein d'un même runner.atick(...), plusieurs types d'événements sont produits ; la fonction _output (_output:4184) ne fait que filtrer/répartir, sans réexécuter le graphe. Passer stream_mode=["values","updates"] ne lance pas deux runs.
  • Namespace unifié des sous-graphes : avec subgraphs=True, les événements portent un préfixe namespace ((ns, mode, payload)) ; le graphe parent voit les événements internes des sous-graphes, le tagging est fait par _output à la sortie de file, plutôt que par le sous-graphe qui ouvrirait son propre flux.
  • v2 typé : via version="v2", l'événement devient un dict {"type": mode, "ns": ..., "data": ..., "interrupts": ...}, plus facile à traiter programmatiquement ; le mode values extrait en plus __interrupt__ dans un champ interrupts séparé.

Fichiers clés

  • Pregel class:450 — définition de la classe Pregel, tous les points d'entrée de streaming s'y trouvent.
  • Pregel.stream:2655 — entrée synchrone du streaming, la signature liste toutes les options stream_mode.
  • Pregel.astream:3063 — entrée async du streaming, la boucle principale async est dans son corps.
  • Pregel.invoke:3836 — entrée de convergence synchrone, fait for chunk in self.stream(...) et prend le dernier chunk values.
  • sync main loop:2964 — boucle principale synchrone : while loop.tick()runner.tickloop.after_tick().
  • async main loop:3437 — boucle principale async : trois phases identiques, runner.atick remplace la version synchrone.
  • _output:4184 — fonction de filtrage à la sortie de file, décide du yield de chaque événement selon stream_mode / print_mode / subgraphs.
  • GraphRunStream:31 — wrapper synchrone piloté par l'appelant, la boucle for est le pump, pas de thread en arrière-plan.
  • AsyncGraphRunStream:304 — pendant async, plusieurs projections (run.values / run.messages) partagent un seul pump.
  • stream_mode attr:709Pregel par défaut à stream_mode="values", peut être surchargé par stream(stream_mode=...).

Flux de données

Le squelette de la boucle principale async est dans le corps de astream (async main loop:3437), trois phases comme décrit dans la page Pregel, la différence est que _output est intercalé entre chaque phase pour vider la file vers l'appelant :

python
while loop.tick():
    for task in await loop.amatch_cached_writes():
        loop.output_writes(task.id, task.writes, cached=True)
    async for _ in runner.atick(
        [t for t in loop.tasks.values() if not t.writes],
        timeout=self.step_timeout,
        get_waiter=get_waiter,
        schedule_task=loop.aaccept_push,
    ):
        # emit output
        for o in _output(
            stream_mode,
            print_mode,
            subgraphs,
            stream.get_nowait,
            asyncio.QueueEmpty,
            version,
            _output_mapper,
            _state_mapper,
        ):
            yield o
    loop.after_tick()
    await aemit_graph_lifecycle_events(loop)
    # wait for checkpoint
    if durability_ == "sync":
        await cast(asyncio.Future, loop._put_checkpoint_fut)

runner.atick réécrit les writes de chaque tâche dans loop, et pousse des triplets (ns, mode, payload) dans stream (asyncio.Queue). Quand la boucle async for externe récupère la main, _output appelle stream.get_nowait() pour vider tous les événements prêts dans la file, filtre selon stream_mode, puis yield à l'utilisateur. La logique centrale de _output (_output:4184) : récupère (ns, mode, payload), si mode in print_mode fait d'abord un print, et si mode in stream_mode seulement yield vers l'extérieur — donc print_mode est un commutateur de debug « afficher sans émettre », il n'affecte pas le contenu réellement yieldé.

Le diagramme suivant traverse l'appel stream de l'entrée jusqu'à la sortie d'événements :

Limites et échecs

  • stream_mode par défaut : quand stream_mode=None, si l'appel est un sous-graphe (CONFIG_KEY_TASK_ID dans le config), défaut à values, sinon utilise self.stream_mode (stream_mode default:2740) — pour éviter que le mode updates par défaut d'un sous-graphe pollue la sortie du graphe parent.
  • Limite de pas : après la sortie de la boucle principale, on regarde loop.status ; si out_of_steps, lève GraphRecursionError en suggérant d'augmenter recursion_limit (out_of_steps:3002) ; l'état draining lève GraphDrained et rend le contrôle au RunControl externe.
  • Sortie multi-mode simultanée : quand stream_mode est une liste, chaque événement devient un tuple (mode, payload) ; avec subgraphs=True en plus, c'est (ns, mode, payload) (tuple mode:4240) — l'appelant doit déstructurer le tuple.
  • Traitement d'héritage du mode messages : quand stream_mode="messages" et version="v1", le v2 messages handler hérité de la couche supérieure est dépouillé, pour éviter que le flux v1 soit routé vers le protocole d'événements en blocs de contenu (strip v2 handler:2773), mais le handler v1 est conservé pour permettre qu'avec subgraphs=True les événements messages internes soient observés par la couche externe.
  • Synchronisation durability et checkpoint : avec durability="sync", la boucle principale await loop._put_checkpoint_fut à chaque étape, garantissant que le checkpoint (point de contrôle) est persisté avant l'étape suivante ; "async" est la valeur par défaut, la persistance se fait en parallèle de l'étape suivante ; "exit" persiste seulement à la sortie.
  • stream_events v3 expérimental : stream_events(version="v3") n'accepte pas les paramètres stream_mode et subgraphs, ceux-ci sont pris en charge internement par le mux, les passer explicitement est rejeté par _reject_v3_invariant_kwargs (_reject_v3_invariant_kwargs:387).

Résumé

astream / stream enveloppent la boucle BSP de Pregel dans une interface itérable, le choix de stream_mode détermine la projection de sortie, et _output est le goulet d'étranglement unique du filtrage en sortie de file. Pour les détails d'exécution, voir /pregel/pregel et /pregel/loop ; pour le mux multi-mode, voir libs/langgraph/langgraph/stream/_mux.py.

Voir la documentation officielle : LangGraph 文档 · README