Stream Transformers: convierten la salida raw de Pregel en proyecciones usables
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 = Trueexponen su proyección como atributo directorun.<key>(_native:50-52); los no nativos van arun.extensionsy 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_modesdeduce qué modos debe emitir el grafo: cada transformer declara qué modos necesita(required_stream_modes:88-93);_collect_stream_modestoma la unión(_collect_stream_modes:398-414), y la llamada v3 alstream()subyacente pasa esa unión comostream_mode(stream_mode=_collect_stream_modes:3549). El usuario no especificastream_modedirectamente, sino que «tira bajo demanda».- Filtrado por scope para evitar ruido:
ValuesTransformer.processcompruebaparams["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 querun.valuesraíz no se contamina con el grafo anidado. - Doble vía sync / async con autodetección: por defecto
supports_sync = False; basta con reescribiraprocess/afinalize/afailpara quetransformer_requires_asynclo 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ícitamentesupports_sync = True, como haceSubgraphTransformer(supports_sync:689). SubgraphTransformerdevuelveSubgraphRunStreamen lugar de eventos raw: al descubrir una invocación de subgrafo,_on_startedconstruye unSubgraphRunStreamque envuelve al mini-mux(SubgraphTransformer:670-705); el usuario no recibe un flujo de eventos, sino un handle sobre el que puede recorrer recursivamentehandle.values/handle.messages/handle.subgraphs, con la misma interfaz isomorfa que elrunraíz.
Archivos clave
StreamTransformer 基类:44-115—init/process/aprocess/finalize/failabstractos + ClassVar de configuración.init:128-138— devuelve el dict de projection; las keys van aextensions, y con_native=Truetambién se montan enrun.<key>.process:148-164— vía sync, por defectoraise NotImplementedError; las subclases deben reescribirla.aprocess:166-184— vía async, por defecto delega enprocess; se reescribe cuando se necesita trabajo async.schedule:233-263—asyncio.Taskligado al ciclo de vida del mux;gatheren aclose,cancelen afail.transformer_requires_async:308-330— detecta siaprocess/afinalize/afailse han reescrito, decidiendo si el transformer puede correr en un mux sync.ValuesTransformer:28-82— proyecciónrun.values, conrequired_stream_modes=("values",).UpdatesTransformer:120-152— proyecciónrun.updates.MessagesTransformer:155-335— proyecciónrun.messages, empaqueta el flujo de tokens del LLM en un handleChatModelStreamy procesa la entrada(payload, metadata)del protocolo v2.LifecycleTransformer:608-667— proyecciónrun.lifecycle, reenvía los eventos de ciclo de vida del subgrafo al registro principal como native protocol event, visible para SDKs remotos.SubgraphTransformer:670-705— proyecciónrun.subgraphs, construye unSubgraphRunStream+ mini-mux al descubrir una invocación de subgrafo.GraphRunStream._pump_next:107-130— pump: extrae un chunk, lo pasa porconvert_to_protocol_event, y lo entrega al mux.convert_to_protocol_event:10-32— mapeo de campos deStreamPartv2 aProtocolEventv3.
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:
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):
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_next → mux.push → transformer.
Límites y fallos
processdebe implementarse al menos uno, einitdebe devolver dict:inites@abstractmethodpor defecto; si una subclase no lo implementa, el arranque falla(init abstract:127-138);processpor defecto lanzaraise NotImplementedError(process default:162-164), y sin reescribirlo estallará al primer evento push.- El orden entre transformers es sensible: los transformers con
before_builtins = Truese ejecutan primero(before_builtins:94-109), pero el docstring advierte de que verán eventostasksantes de que los consumanLifecycleTransformer/SubgraphTransformer; modificarevent["params"]["namespace"]o los camposid/result/error/interruptsdescuadrará la contabilidad de los dos siguientes, por lo que sólo se permite observar o modificar campos fríos. MessagesTransformerestá acoplado al protocolo v2:stream_events(version="v3")fuerzaversion="v2"yCONFIG_KEY_STREAM_MESSAGES_V2=Trueal invocar elstream()subyacente(_pregel_stream_v3 stream call:3545-3551); de lo contrario, la proyección messages no obtiene el formato correcto de payload.StreamChanneles de un único consumidor; iterar dos veces lanza error:run.valuessólo puede ser consumido por un buclefor(single-consumer note:38-40); para hacer fan-out hay que usarprojection.tee(n)(tee:245), con el coste de un buffer.schedule(on_error="raise")convierte el cierre en fail:on_error="raise"enschedulehace que una excepción del task empuje al mux aafailduranteaclose(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_trackestrictamente 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 elhandle.subgraphsdel handle hijo; si el usuario saltahandle.subgraphsy sólo mirarun.subgraphsraí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