Skip to content

Stream Transformers: convierten la salida raw de Pregel en proyecciones usables

源码版本1.2.9

Responsabilidades

PregelLoop, en cada superstep, utiliza _emit para introducir tuplas (checkpoint_ns, mode, payload) en un SyncQueue(_emit:1380-1414), con eventos raw de tipo «valor», «actualización», «task», «checkpoint», etc. Ésta es la forma de salida de los protocolos v1 / v2 — el usuario recibe una secuencia de chunks. Pero en el protocolo v3, el usuario no quiere chunks, sino «proyecciones (projection) organizadas por modo, iterables por un único consumidor»: run.values es un flujo de instantáneas completas del estado, run.messages es un conjunto de handles ChatModelStream para flujos de tokens del LLM, run.updates es un flujo de actualizaciones de nodos, y run.subgraphs es un handle para entrar en el subgrafo.

La lógica que convierte los chunks raw en estas proyecciones se encuentra en libs/langgraph/langgraph/stream/transformers.py(ValuesTransformer:28), un transformer por modo, todos heredando de la clase base abstracta StreamTransformer(StreamTransformer:44). Son registrados por StreamMux, reciben eventos de GraphRunStream, y mantienen internamente un StreamChannel como proyección. Al invocar graph.stream_events(version="v3"), el usuario obtiene en realidad un GraphRunStream(GraphRunStream:31), cuyos atributos .values / .messages / .updates / .subgraphs / .custom / .lifecycle son las entradas de proyección de los transformers nativos(native attrs:78-79).

Todo el sistema de transformers es el punto de extensión del protocolo v3 — además de los cinco internos, el usuario puede añadir los suyos en Pregel.compile(stream_transformers=[...])(stream_transformers param:783), que en _pregel_stream_v3 se pasan junto a los transformers internos al StreamMux en forma de factory(factories list:3533-3544).

Motivación de diseño

  • Doble vía nativa vs extensión: los transformers con _native = True exponen su proyección como atributo directo run.<key>(_native:50-52); los no nativos van a run.extensions y sólo se exponen en el dict. Los cinco internos son todos nativos; los transformers personalizados entran por defecto en extensions, y para exponerlos como atributo directo hay que fijar explícitamente _native = True.
  • required_stream_modes deduce qué modos debe emitir el grafo: cada transformer declara qué modos necesita(required_stream_modes:88-93); _collect_stream_modes toma la unión(_collect_stream_modes:398-414), y la llamada v3 al stream() subyacente pasa esa unión como stream_mode(stream_mode=_collect_stream_modes:3549). El usuario no especifica stream_mode directamente, sino que «tira bajo demanda».
  • Filtrado por scope para evitar ruido: ValuesTransformer.process comprueba params["namespace"] != self._scope_list(ValuesTransformer.process:70-82), aceptando sólo eventos de su scope; los eventos del subgrafo los maneja el mini-mux del propio subgrafo, de modo que run.values raíz no se contamina con el grafo anidado.
  • Doble vía sync / async con autodetección: por defecto supports_sync = False; basta con reescribir aprocess / afinalize / afail para que transformer_requires_async lo marque como async-only(transformer_requires_async:308-330), y el registro en un mux síncrono falla al instante. Para soportar sync y async a la vez hay que fijar explícitamente supports_sync = True, como hace SubgraphTransformer(supports_sync:689).
  • SubgraphTransformer devuelve SubgraphRunStream en lugar de eventos raw: al descubrir una invocación de subgrafo, _on_started construye un SubgraphRunStream que envuelve al mini-mux(SubgraphTransformer:670-705); el usuario no recibe un flujo de eventos, sino un handle sobre el que puede recorrer recursivamente handle.values / handle.messages / handle.subgraphs, con la misma interfaz isomorfa que el run raíz.

Archivos clave

Flujo de datos

Toda la cadena v3 comienza en PregelLoop._emit, pasa por SyncQueue, GraphRunStream._pump_next, convert_to_protocol_event, y termina en el mux. convert_to_protocol_event convierte el {type, ns, data, interrupts} de v2 en un ProtocolEvent de 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) Obsérvese que method es el type de v2 — values / updates / messages / custom / tasks / checkpoints / debug — y es la clave usada por StreamTransformer.process para dispatch; por ejemplo, ValuesTransformer sólo actúa sobre eventos con method == "values"(method check:70-71).

En cambio, en el protocolo v1 / v2, Pregel.stream recorre otra ruta: los chunks raw salen directamente del SyncQueue y pasan por la función _output(_output:4184-4243), que según si stream_mode es una cadena o una lista decide devolver payload, (mode, payload) o (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) Por eso stream_mode=["updates", "values"] devuelve tuplas de tres elementos mientras que una cadena única devuelve un valor único — un mismo chunk raw se formatea de distinta manera en _output. v3 no pasa por _output, sino por GraphRunStream._pump_nextmux.push → transformer.

Límites y fallos

  • process debe implementarse al menos uno, e init debe devolver dict: init es @abstractmethod por defecto; si una subclase no lo implementa, el arranque falla(init abstract:127-138); process por defecto lanza raise NotImplementedError(process default:162-164), y sin reescribirlo estallará al primer evento push.
  • El orden entre transformers es sensible: los transformers con before_builtins = True se ejecutan primero(before_builtins:94-109), pero el docstring advierte de que verán eventos tasks antes de que los consuman LifecycleTransformer / SubgraphTransformer; modificar event["params"]["namespace"] o los campos id / result / error / interrupts descuadrará la contabilidad de los dos siguientes, por lo que sólo se permite observar o modificar campos fríos.
  • MessagesTransformer está acoplado al protocolo v2: stream_events(version="v3") fuerza version="v2" y CONFIG_KEY_STREAM_MESSAGES_V2=True al invocar el stream() subyacente(_pregel_stream_v3 stream call:3545-3551); de lo contrario, la proyección messages no obtiene el formato correcto de payload.
  • StreamChannel es de un único consumidor; iterar dos veces lanza error: run.values sólo puede ser consumido por un bucle for(single-consumer note:38-40); para hacer fan-out hay que usar projection.tee(n)(tee:245), con el coste de un buffer.
  • schedule(on_error="raise") convierte el cierre en fail: on_error="raise" en schedule hace que una excepción del task empuje al mux a afail durante aclose(on_error:254-257); el valor por defecto "log" sólo registra sin propagar — elegir mal puede hacer que el fallo de un único transformer arrastre toda la corriente.
  • SubgraphTransformer._should_track estrictamente mayor que el scope: sólo rastrea subgrafos del nivel inmediatamente inferior al scope actual(should_track:634-636); los subgrafos nietos se descubren recursivamente vía el handle.subgraphs del handle hijo; si el usuario salta handle.subgraphs y sólo mira run.subgraphs raíz, los eventos del subgrafo nieto no aparecerán ahí.

Resumen

StreamTransformer es el punto de extensión del protocolo de streaming v3: cinco transformers internos convierten los raw (ns, mode, payload) en proyecciones tipadas iterables, y el usuario puede añadir los suyos propios vía el parámetro stream_transformers. Toda la maquinaria de dispatch se describe en StreamMux; para ver cómo los chunks raw entran al SyncQueue desde PregelLoop, ver PregelLoop.

Véase la documentación oficial: Documentación de LangGraph · README