Skip to content

Salida en streaming: astream y stream_mode

源码版本1.2.9

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 / atick inserta eventos conforme completa cada task, y el bucle for externo 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. Pasar stream_mode=["values","updates"] no ejecuta el grafo dos veces.
  • Namespace unificado para subgrafos: con subgraphs=True activado, los eventos llevan un prefijo namespace ((ns, mode, payload)); el grafo padre ve los eventos internos del subgrafo, gracias a que _output los 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 modo values extrae __interrupt__ a un campo interrupts independiente.

Archivos clave

  • Pregel class:450 — definición de la clase Pregel; todas las entradas de streaming están aquí.
  • Pregel.stream:2655 — entrada de streaming síncrono; en la firma se listan todas las opciones de stream_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; internamente for chunk in self.stream(...) tomando el último chunk values.
  • sync main loop:2964 — bucle principal síncrono: while loop.tick()runner.tickloop.after_tick().
  • async main loop:3437 — bucle principal asíncrono: tres etapas como arriba, con runner.atick en 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ún stream_mode / print_mode / subgraphs.
  • GraphRunStream:31 — wrapper síncrono de stream dirigido por el llamador; el bucle for es 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:709Pregel tiene por defecto stream_mode="values", que puede sobrescribirse con stream(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:

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)

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_mode por defecto: cuando stream_mode=None, si se invoca como subgrafo (CONFIG_KEY_TASK_ID en config), el valor por defecto es values; en caso contrario se usa self.stream_mode(stream_mode default:2740), evitando que el modo updates por defecto del subgrafo contamine la salida del padre.
  • Límite de pasos: al salir del bucle principal se observa loop.status; si es out_of_steps se lanza GraphRecursionError, sugiriendo subir recursion_limit(out_of_steps:3002); el estado draining lanza GraphDrained, devolviendo el control al RunControl externo.
  • Salida multimodo simultánea: cuando stream_mode se pasa como lista, cada evento se convierte en una tupla (mode, payload); con subgraphs=True ademá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" y version="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 eventos messages internos sean observados por la capa exterior cuando subgraphs=True.
  • durability y sincronización de checkpoint: con durability="sync", el bucle principal hace await loop._put_checkpoint_fut despué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ámetros stream_mode ni subgraphs, 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