Send : primitif de tâche PUSH pour le fan-out de nœuds
Responsabilités
Send est le primitif de LangGraph pour « livrer une tâche à un nœud donné », défini dans libs/langgraph/langgraph/types.py(Send class:664). Un Send(node="foo", arg={"k": v}) signifie « au prochain pas, appeler le nœud foo avec arg comme entrée » ; la différence avec un déclenchement de nœud normal tient à ce qu'un nœud normal (PULL) est déclenché passivement par la relation d'abonnement aux canaux, tandis que Send (PUSH) est déclaré activement par le nœud précédent : « je veux que foo s'exécute une fois, avec cet arg spécifique ».
Sa position se situe entre les conditional edges et le canal TASKS : add_conditional_edges(START, lambda s: [Send("gen", {"x": 1}), Send("gen", {"x": 2})]) écrit la liste de Send renvoyée par la fonction conditionnelle dans le canal TASKS(goto as sends:62), et au tour suivant prepare_next_tasks les récupère via tasks_channel.get(), en assemblant chaque Send en un PregelExecutableTask indépendant(PUSH tasks:421).
Motivation de conception
- Expression naturelle de map-reduce : le scénario classique est « exécuter N fois le même nœud en parallèle, avec une entrée différente à chaque fois ».
Send("generate", {"subject": s})renvoyé dans une liste par un conditional edge lance chaque Send comme un appel indépendant ; Pregel les exécute naturellement en concurrence, et le reducer fusionne les N résultats dans l'état principal. Pas besoin de dérouler manuellement N arêtes dans le graphe. - Dichotomie PUSH / PULL : le déclenchement d'un nœud est de deux types — PULL (passif, déclenché par la mise à jour d'un canal souscrit) et PUSH (actif, délivré explicitement par
Send).prepare_next_tasksconsomme d'abord les tâches PUSH du canalTASKS, puis calcule les tâches PULL viatrigger_to_nodes(PUSH tasks:440)`. argn'a pas à être l'état principal :Send.argpeut être n'importe quel objet, sans devoir respecter le schéma d'état principal du graphe — c'est l'input_schemadu nœud cible qui détermine la forme dearg, ce qui permet de « livrer des entrées hétérogènes au même nœud », par ex.Send("review", ReviewContext(doc=...))vsSend("review", DefaultContext()).Sendéquivalent àCommand.goto:Command(goto=Send(...))passe par le même canalTASKSdansmap_command(goto as sends:62), donc conditional edges et valeur de retour de nœud convergent vers l'ordonnancement PUSH, abstraction unifiée.- Chaque
Sendest une task indépendante : N éléments dans une liste deSenddeviennent NPregelExecutableTaskindépendants, avec leur propre task_id et path, pouvant être retentés et gérer leurs erreurs séparément — l'échec de l'un n'affecte pas les autres Send concurrents.
Fichiers clés
Send class:664— dataclassSend(node, arg, *, timeout=None),__slots__ = ("node", "arg", "timeout").Send __init__:718— constructeur,timeoutpasse parTimeoutPolicy.coerce.Send __hash__:739—hash((node, arg, timeout)), utilisé comme clé de déduplication dans le dictfutures.prepare_next_tasks:392— fonction d'ordonnancement principale, consomme d'abord les tâches PUSH puis calcule les nœuds candidats PULL.PUSH tasks:440— récupère la séquence de Send depuis le canalTASKS, et appelleprepare_single_taskpour assembler chaquePregelExecutableTask.PULL tasks:470— chemin PULL, calcule quels nœuds doivent s'exécuter selontrigger_to_nodes, ne consomme pasTASKS.goto as sends:62—map_commandécrit lesSenddeCommand.gotodans le canalTASKS, unifiant le chemin PUSH.add_conditional_edges:969— entrée des conditional edges, la fonction conditionnelle renvoie unSendou une liste deSend.PregelExecutableTask:627— forme exécutable de la tâche, détientname/input/proc/writes/triggers/pathetc.PregelRunner.tick:176— exécute en concurrence toutes les tâches PUSH + PULL de cette étape ; les tâches issues deSendy sont exécutées.
Flux de données
Send est réellement consommé au début de prepare_next_tasks(prepare_next_tasks:392) :
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)Le canal TASKS est un Topic[Send] — un type de canal qui supporte plusieurs écritures et plusieurs lectures(Topic:23). Chaque Send est récupéré dans tasks_channel.get(), et (PUSH, idx) sert de préfixe de task_id passé à prepare_single_task, qui récupère la configuration PregelNode correspondante depuis processes[send.node] (incluant proc / writers / retry_policy etc.), utilise send.arg comme input pour assembler un PregelExecutableTask. Send.arg ne passe pas par le canal d'état principal — il est passé directement comme task.input au proc du nœud ; la fonction de nœud reçoit cet arg, pas un snapshot complet de l'état principal du graphe.
Send entre dans le canal TASKS par deux chemins : (1) la fonction conditionnelle configurée par add_conditional_edges renvoie une liste de Send, et attach_branch les écrit dans TASKS ; (2) la fonction de nœud renvoie Command(goto=Send(...)), et map_command dans _io.py écrit le Send dans TASKS(goto as sends:62). Les deux chemins convergent : écriture dans le canal TASKS, puis consommation par prepare_next_tasks au tour suivant.
Limites et échecs
Senddoit cibler un nœud enregistré : siprocesses[send.node]n'existe pas,prepare_single_taskrenvoieNoneet le Send est silencieusement abandonné — pas de raise, mais la tâche ne s'exécute pas. En debug, si un Send ne prend pas effet, vérifier d'abord que le nom du nœud cible a bien été enregistré viaadd_node.Send.argne passe pas par le reducer d'état : l'argdeSendsert directement d'input au nœud, sans passer par le reducer du canal d'état principal. Si le nœud cible s'attend à lire l'état principal, il doit le récupérer dans sa fonction viaconfig[CONFIG_KEY_READ]ou un paramètrestate: State—argn'est qu'une « entrée supplémentaire ».- Déduplication de liste de
Sendpar hash :Send.__hash__utilise le triplet(node, arg, timeout)(Send __hash__:739). Si deux Send vers le même nœud ont unargnon hachable (comme un dict), cela lèveTypeError—Sendest utilisé comme clé de dict dans le chemin de déduplication de_call. - Un même nœud peut être déclenché à la fois par PULL et PUSH : si
Send("foo", arg)livre àfooalors qu'un canal trigger defooest aussi mis à jour,prepare_next_tasksliste à la fois la tâche PUSH et la tâche PULL, etPregelRunner.tickles lance en concurrence — c'est autorisé, mais attention que la fonction de nœud peut être appelée deux fois au cours d'une même étape. Send.timeouts'applique seule :Send(node, arg, timeout=10)ne s'applique qu'à cette tâche PUSH, sans affecter la configurationtimeoutpar défaut du nœud cible(Send __init__:718).- Le canal
TASKSest unTopic, pas unLastValue:Topicsupporte l'accumulation de plusieurs écritures(Topic:23), donc plusieurs nœuds peuvent écrire desSenddans le même superpas, et le tour suivant les consomme tous d'un coup — c'est la base de map-reduce.
Résumé
Send est le support des tâches PUSH, en complément des nœuds PULL normaux. Les deux entrées (conditional edges et Command.goto) convergent vers le même canal TASKS, ordonnancé de façon unifiée par prepare_next_tasks. Pour suivre comment les tasks sont exécutées en concurrence voir /func/concurrency, pour la relation conditional edge → nœud voir libs/langgraph/langgraph/graph/_branch.py. Voir la documentation officielle : LangGraph 文档 · README