Skip to content

Send: primitiva de task PUSH para fan-out de nodos

源码版本1.2.9

Responsabilidades

Send es la primitiva de LangGraph para «entregar un task a un nodo específico», definida en libs/langgraph/langgraph/types.py(Send class:664). Un Send(node="foo", arg={"k": v}) indica «en el siguiente paso, invoca el nodo foo con arg como entrada»; la diferencia con la activación de un nodo normal radica en que un nodo normal (PULL) se dispara pasivamente por la relación de suscripción a canales, mientras que Send (PUSH) es una declaración activa del nodo anterior: «quiero que foo se ejecute una vez, con este arg concreto».

Su posición está entre las conditional edges y el canal TASKS: add_conditional_edges(START, lambda s: [Send("gen", {"x": 1}), Send("gen", {"x": 2})]) escribe la lista de Send devuelta por la función condicional en el canal TASKS(goto as sends:62), y en la siguiente ronda prepare_next_tasks los saca de tasks_channel.get(), ensamblando cada Send en un PregelExecutableTask independiente(PUSH tasks:421).

Motivación de diseño

  • Expresión natural de map-reduce: el escenario clásico es «el mismo nodo ejecutándose N veces en paralelo con entradas diferentes». Send("generate", {"subject": s}) devuelve una lista en la conditional edge, cada Send es una invocación independiente, Pregel las ejecuta concurrentemente de forma natural y el reducer fusiona los N resultados de vuelta al estado principal. No hace falta expandir manualmente N aristas en el grafo.
  • Dicotomía PUSH / PULL: la activación de un nodo se divide en dos tipos — PULL (pasiva, disparada por actualizaciones de canales suscritos) y PUSH (activa, entregada explícitamente por Send). prepare_next_tasks consume primero los tasks PUSH del canal TASKS, y luego calcula los tasks PULL según trigger_to_nodes(PUSH tasks:440`).
  • arg no tiene que ser el estado principal: Send.arg puede ser cualquier objeto, sin obligación de cumplir el schema del estado principal del grafo — el input_schema del nodo destino determina la forma de arg, permitiendo «entregar entradas heterogéneas al mismo nodo», por ejemplo Send("review", ReviewContext(doc=...)) vs Send("review", DefaultContext()).
  • Send equivale a Command.goto: Command(goto=Send(...)) pasa por el mismo canal TASKS dentro de map_command(goto as sends:62), de modo que las dos entradas — conditional edges y valores de retorno de nodos — confluyen finalmente en la programación de tasks PUSH, con un único nivel de abstracción.
  • Cada Send es un task independiente: los N elementos de una lista de Send son N PregelExecutableTask independientes, cada uno con su propio task_id y path, que pueden reintentrarse y tratarse ante errores por separado — el fallo de uno no afecta a los demás Sends concurrentes.

Archivos clave

  • Send class:664 — dataclass Send(node, arg, *, timeout=None), con __slots__ = ("node", "arg", "timeout").
  • Send __init__:718 — constructor; timeout pasa por TimeoutPolicy.coerce.
  • Send __hash__:739hash((node, arg, timeout)), base de deduplicación como key en el dict futures.
  • prepare_next_tasks:392 — función principal de programación; consume primero tasks PUSH y luego calcula los nodos candidatos PULL.
  • PUSH tasks:440 — saca la secuencia de Send del canal TASKS, llamando a prepare_single_task para cada uno y ensamblando un PregelExecutableTask.
  • PULL tasks:470 — ruta PULL, calcula qué nodos deben ejecutarse según trigger_to_nodes, sin consumir TASKS.
  • goto as sends:62map_command escribe el Send de Command.goto en el canal TASKS, unificando la ruta PUSH.
  • add_conditional_edges:969 — entrada de conditional edge; la función condicional devuelve un Send o una lista de Send.
  • PregelExecutableTask:627 — forma ejecutable del task, contiene name / input / proc / writes / triggers / path, etc.
  • PregelRunner.tick:176 — ejecuta concurrentemente todos los tasks PUSH + PULL de este paso; aquí se ejecutan los tasks entregados por Send.

Flujo de datos

Send se consume realmente en prepare_next_tasks(prepare_next_tasks:392), al inicio:

python
input_cache: dict[INPUT_CACHE_KEY_TYPE, Any] = {}
checkpoint_id_bytes = binascii.unhexlify(checkpoint["id"].replace("-", ""))
null_version = checkpoint_null_version(checkpoint)
tasks: list[PregelTask | PregelExecutableTask] = []
# Consume pending tasks
tasks_channel = cast(Topic[Send] | None, channels.get(TASKS))
if tasks_channel and tasks_channel.is_available():
    for idx, _ in enumerate(tasks_channel.get()):
        if task := prepare_single_task(
            (PUSH, idx),
            None,
            checkpoint=checkpoint,
            checkpoint_id_bytes=checkpoint_id_bytes,
            checkpoint_null_version=null_version,
            pending_writes=pending_writes,
            processes=processes,
            channels=channels,
            managed=managed,
            config=config,
            step=step,
            stop=stop,
            for_execution=for_execution,
            store=store,
            checkpointer=checkpointer,
            manager=manager,
            input_cache=input_cache,
            cache_policy=cache_policy,
            retry_policy=retry_policy,
        ):
            tasks.append(task)

El canal TASKS es un Topic[Send] — un tipo de canal que soporta escritura y lectura múltiples(Topic:23). Cada Send se extrae en tasks_channel.get(), con (PUSH, idx) como prefijo de task_id pasado a prepare_single_task, que a su vez obtiene de processes[send.node] la configuración PregelNode correspondiente (incluyendo proc / writers / retry_policy, etc.) y ensambla send.arg como input en un PregelExecutableTask. Send.arg no entra en el canal de estado principal — se pasa directamente como task.input al proc del nodo; la función nodo recibe este arg, no la instantánea completa del estado principal del grafo.

Send entra al canal TASKS por dos rutas: (1) la función condicional configurada por add_conditional_edges devuelve una lista de Send, y attach_branch los escribe en TASKS; (2) la función nodo devuelve Command(goto=Send(...)), y map_command en _io.py escribe el Send en TASKS(goto as sends:62). Ambas rutas confluyen en el mismo canal TASKS, y la siguiente ronda de prepare_next_tasks los consume.

Límites y fallos

  • Send debe apuntar a un nodo registrado: si processes[send.node] no existe, prepare_single_task devuelve None y el Send se descarta silenciosamente — no se lanza ningún error, pero el task no se ejecuta. Al depurar, si un Send no hace efecto, comprueba primero si el nombre del nodo destino se registró con add_node.
  • Send.arg no pasa por el reducer de estado: el arg de Send se pasa directamente como input del nodo, sin fusionarse mediante el reducer del canal de estado principal. Si el nodo destino espera leer el estado principal, debe obtenerlo dentro de su función a través de config[CONFIG_KEY_READ] o del parámetro state: Statearg es sólo una «entrada adicional».
  • La deduplicación de listas de Send se basa en hash: Send.__hash__ usa la tupla (node, arg, timeout)(Send __hash__:739); si los arg de dos Send al mismo nodo son objetos no hashables (como un dict), se lanza TypeError directamente — Send se usa como key de dict en la ruta de deduplicación en _call.
  • Un mismo nodo puede ser activado por PULL y PUSH simultáneamente: si Send("foo", arg) entrega a foo y a la vez algún canal trigger de foo se actualiza, prepare_next_tasks listará tanto el task PUSH como el PULL, y PregelRunner.tick los ejecutará concurrentemente — esto está permitido, pero ten en cuenta que la función nodo puede invocarse dos veces en el mismo paso.
  • Send.timeout actúa de forma aislada: Send(node, arg, timeout=10) sólo aplica a ese task PUSH, sin afectar a la configuración timeout por defecto del nodo destino(Send __init__:718).
  • El canal TASKS es Topic, no LastValue: Topic soporta la acumulación de múltiples escrituras(Topic:23), por lo que varios nodos pueden escribir Send en un mismo superstep y la siguiente ronda los consume todos de una vez — ésta es la base de map-reduce.

Resumen

Send es el portador de los tasks PUSH, complementario a los nodos PULL normales. Las dos entradas — conditional edges y Command.goto — confluyen en el mismo canal TASKS, programadas de forma unificada por prepare_next_tasks. Para seguir cómo se ejecutan los tasks concurrentemente ver /func/concurrency; para la relación entre conditional edges y nodos ver libs/langgraph/langgraph/graph/_branch.py. Véase la documentación oficial: Documentación de LangGraph · README