Send: primitiva de task PUSH para fan-out de nodos
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_tasksconsume primero los tasks PUSH del canalTASKS, y luego calcula los tasks PULL segúntrigger_to_nodes(PUSH tasks:440`). argno tiene que ser el estado principal:Send.argpuede ser cualquier objeto, sin obligación de cumplir el schema del estado principal del grafo — elinput_schemadel nodo destino determina la forma dearg, permitiendo «entregar entradas heterogéneas al mismo nodo», por ejemploSend("review", ReviewContext(doc=...))vsSend("review", DefaultContext()).Sendequivale aCommand.goto:Command(goto=Send(...))pasa por el mismo canalTASKSdentro demap_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
Sendes un task independiente: los N elementos de una lista deSendson NPregelExecutableTaskindependientes, 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— dataclassSend(node, arg, *, timeout=None), con__slots__ = ("node", "arg", "timeout").Send __init__:718— constructor;timeoutpasa porTimeoutPolicy.coerce.Send __hash__:739—hash((node, arg, timeout)), base de deduplicación como key en el dictfutures.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 canalTASKS, llamando aprepare_single_taskpara cada uno y ensamblando unPregelExecutableTask.PULL tasks:470— ruta PULL, calcula qué nodos deben ejecutarse segúntrigger_to_nodes, sin consumirTASKS.goto as sends:62—map_commandescribe elSenddeCommand.gotoen el canalTASKS, unificando la ruta PUSH.add_conditional_edges:969— entrada de conditional edge; la función condicional devuelve unSendo una lista deSend.PregelExecutableTask:627— forma ejecutable del task, contienename/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 porSend.
Flujo de datos
Send se consume realmente en prepare_next_tasks(prepare_next_tasks:392), al inicio:
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
Senddebe apuntar a un nodo registrado: siprocesses[send.node]no existe,prepare_single_taskdevuelveNoney 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ó conadd_node.Send.argno pasa por el reducer de estado: elargdeSendse 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 deconfig[CONFIG_KEY_READ]o del parámetrostate: State—arges sólo una «entrada adicional».- La deduplicación de listas de
Sendse basa en hash:Send.__hash__usa la tupla(node, arg, timeout)(Send __hash__:739); si losargde dos Send al mismo nodo son objetos no hashables (como un dict), se lanzaTypeErrordirectamente —Sendse 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 afooy a la vez algún canal trigger defoose actualiza,prepare_next_taskslistará tanto el task PUSH como el PULL, yPregelRunner.ticklos ejecutará concurrentemente — esto está permitido, pero ten en cuenta que la función nodo puede invocarse dos veces en el mismo paso. Send.timeoutactúa de forma aislada:Send(node, arg, timeout=10)sólo aplica a ese task PUSH, sin afectar a la configuracióntimeoutpor defecto del nodo destino(Send __init__:718).- El canal
TASKSesTopic, noLastValue:Topicsoporta la acumulación de múltiples escrituras(Topic:23), por lo que varios nodos pueden escribirSenden 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