Salida en streaming: astream y stream_mode
Responsabilidades
Pregel.astream / Pregel.stream son las entradas de ejecución que el grafo compilado expone hacia el exterior. Tú le pasas una entrada, y arranca PregelLoop que va iterando ronda a ronda; durante el proceso, proyecta los datos producidos en cada superstep (superstep) según el stream_mode especificado por el llamador, los coloca en una cola y el llamador los consume por iteración. ainvoke / invoke son esencialmente astream / stream convergiendo a un único valor — el cuerpo de implementación de Pregel.invoke es for chunk in self.stream(...) tomando el último fragmento(Pregel.invoke body:3891).
La responsabilidad de esta capa es desacoplar «ejecución» y «forma de salida»: la ejecución sigue siendo el bucle BSP de Pregel, y la forma de salida la determina stream_mode — para una misma ejecución del grafo, puedes pedir values (estado completo en cada paso), updates (incrementales en cada paso), messages (flujo de tokens del LLM), custom (salida personalizada dentro del nodo), checkpoints (eventos de checkpoint), tasks (inicio y fin de tasks), debug (todo), etc., e incluso pasar una lista y obtener varios modos a la vez.
Motivación de diseño
- No bloquear la ejecución: la salida va por una cola;
runner.tick/atickinserta eventos conforme completa cada task, y el bucleforexterno itera y recibe a medida, en lugar de esperar a que termine todo el grafo para devolverlo de golpe. Esto es imprescindible para el flujo de tokens del LLM (messages) — el usuario necesita verlos según se generan. - Múltiples proyecciones reutilizan la misma ejecución: dentro del mismo
runner.atick(...)se generan varios tipos de eventos; la función_output(_output:4184) sólo filtra o reparte, no vuelve a ejecutar el grafo. Pasarstream_mode=["values","updates"]no ejecuta el grafo dos veces. - Namespace unificado para subgrafos: con
subgraphs=Trueactivado, los eventos llevan un prefijonamespace((ns, mode, payload)); el grafo padre ve los eventos internos del subgrafo, gracias a que_outputlos etiqueta al sacarlos de la cola, sin que el subgrafo tenga que abrir su propio stream. - v2 tipado: la nueva versión, mediante
version="v2", convierte el evento en un dict{"type": mode, "ns": ..., "data": ..., "interrupts": ...}, facilitando el procesamiento programático posterior; además, el modovaluesextrae__interrupt__a un campointerruptsindependiente.
Archivos clave
Pregel class:450— definición de la clasePregel; todas las entradas de streaming están aquí.Pregel.stream:2655— entrada de streaming síncrono; en la firma se listan todas las opciones destream_mode.Pregel.astream:3063— entrada de streaming asíncrono; el bucle principal asíncrono se encuentra en su cuerpo de implementación.Pregel.invoke:3836— entrada síncrona convergente; internamentefor chunk in self.stream(...)tomando el último chunkvalues.sync main loop:2964— bucle principal síncrono:while loop.tick()→runner.tick→loop.after_tick().async main loop:3437— bucle principal asíncrono: tres etapas como arriba, conrunner.aticken lugar de la versión síncrona._output:4184— función de filtrado al sacar de la cola, que decide cómo yield cada evento segúnstream_mode/print_mode/subgraphs.GraphRunStream:31— wrapper síncrono de stream dirigido por el llamador; el buclefores el pump, sin hilos en background.AsyncGraphRunStream:304— contrapartida asíncrona, donde varias proyecciones (run.values/run.messages) comparten un único pump.stream_mode attr:709—Pregeltiene por defectostream_mode="values", que puede sobrescribirse constream(stream_mode=...).
Flujo de datos
El esqueleto del bucle principal asíncrono está en el cuerpo de astream(async main loop:3437), con la misma estructura de tres etapas que en la página del motor Pregel, con la diferencia de que entre cada etapa se intercala _output para entregar al llamador lo que hay en la cola:
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)Dentro de runner.atick, los writes generados por cada task se escriben en loop, y simultáneamente las tuplas (ns, mode, payload) se insertan en stream (un asyncio.Queue); cuando el bucle async for externo recibe el control, _output llama a stream.get_nowait() para vaciar todos los eventos listos en la cola, los filtra según stream_mode y los yield al usuario. La lógica central de _output(_output:4184): obtiene (ns, mode, payload), si mode in print_mode primero hace print, y sólo hace yield hacia fuera si mode in stream_mode — por tanto print_mode es un interruptor de depuración que «sólo observa, no envía», sin afectar al contenido realmente yieltado.
El siguiente diagrama recorre la llamada a stream desde la entrada hasta que el evento sale de la cola:
Límites y fallos
stream_modepor defecto: cuandostream_mode=None, si se invoca como subgrafo (CONFIG_KEY_TASK_IDen config), el valor por defecto esvalues; en caso contrario se usaself.stream_mode(stream_mode default:2740), evitando que el modoupdatespor defecto del subgrafo contamine la salida del padre.- Límite de pasos: al salir del bucle principal se observa
loop.status; si esout_of_stepsse lanzaGraphRecursionError, sugiriendo subirrecursion_limit(out_of_steps:3002); el estadodraininglanzaGraphDrained, devolviendo el control alRunControlexterno. - Salida multimodo simultánea: cuando
stream_modese pasa como lista, cada evento se convierte en una tupla(mode, payload); consubgraphs=Trueademás se convierte en(ns, mode, payload)(tuple mode:4240), y el llamador debe deconstruir la tupla. - Tratamiento heredado del modo messages: con
stream_mode="messages"yversion="v1", se descarta el messages handler v2 heredado de capas superiores, evitando que el flujo v1 se enrute al protocolo de eventos por content-block(strip v2 handler:2773), pero conservando el handler v1 para soportar que los eventosmessagesinternos sean observados por la capa exterior cuandosubgraphs=True. - durability y sincronización de checkpoint: con
durability="sync", el bucle principal haceawait loop._put_checkpoint_futdespués de cada paso, garantizando que el checkpoint se persista antes de avanzar al siguiente;"async"(valor por defecto) persiste en paralelo con el siguiente paso;"exit"sólo persiste al salir. - stream_events v3 experimental:
stream_events(version="v3")no acepta los parámetrosstream_modenisubgraphs, que son gestionados internamente por el mux; pasarlos explícitamente es rechazado por_reject_v3_invariant_kwargs(_reject_v3_invariant_kwargs:387).
Resumen
astream / stream envuelven el bucle BSP de Pregel en una interfaz iterable; la elección de stream_mode determina la proyección de salida, y _output es el único embudo de filtrado al sacar de la cola. Para detalles de ejecución ver /pregel/pregel y /pregel/loop; para el mux multimodo ver libs/langgraph/langgraph/stream/_mux.py. Véase la documentación oficial: Documentación de LangGraph · README