Streaming de sortie : astream et stream_mode
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/aticky pousse un événement dès qu'une tâche est terminée, la boucleforexterne 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. Passerstream_mode=["values","updates"]ne lance pas deux runs. - Namespace unifié des sous-graphes : avec
subgraphs=True, les événements portent un préfixenamespace((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 modevaluesextrait en plus__interrupt__dans un champinterruptsséparé.
Fichiers clés
Pregel class:450— définition de la classePregel, tous les points d'entrée de streaming s'y trouvent.Pregel.stream:2655— entrée synchrone du streaming, la signature liste toutes les optionsstream_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, faitfor chunk in self.stream(...)et prend le dernier chunkvalues.sync main loop:2964— boucle principale synchrone :while loop.tick()→runner.tick→loop.after_tick().async main loop:3437— boucle principale async : trois phases identiques,runner.atickremplace la version synchrone._output:4184— fonction de filtrage à la sortie de file, décide du yield de chaque événement selonstream_mode/print_mode/subgraphs.GraphRunStream:31— wrapper synchrone piloté par l'appelant, la boucleforest 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:709—Pregelpar défaut àstream_mode="values", peut être surchargé parstream(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 :
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_IDdans le config), défaut àvalues, sinon utiliseself.stream_mode(stream_mode default:2740) — pour éviter que le modeupdatespar 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; siout_of_steps, lèveGraphRecursionErroren suggérant d'augmenterrecursion_limit(out_of_steps:3002) ; l'étatdraininglèveGraphDrainedet rend le contrôle auRunControlexterne. - Sortie multi-mode simultanée : quand
stream_modeest une liste, chaque événement devient un tuple(mode, payload); avecsubgraphs=Trueen 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"etversion="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'avecsubgraphs=Trueles événementsmessagesinternes soient observés par la couche externe. - Synchronisation durability et checkpoint : avec
durability="sync", la boucle principaleawait 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_eventsv3 expérimental :stream_events(version="v3")n'accepte pas les paramètresstream_modeetsubgraphs, 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