Skip to content

Send : primitif de tâche PUSH pour le fan-out de nœuds

源码版本1.2.9

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_tasks consomme d'abord les tâches PUSH du canal TASKS, puis calcule les tâches PULL via trigger_to_nodes(PUSH tasks:440)`.
  • arg n'a pas à être l'état principal : Send.arg peut être n'importe quel objet, sans devoir respecter le schéma d'état principal du graphe — c'est l'input_schema du nœud cible qui détermine la forme de arg, ce qui permet de « livrer des entrées hétérogènes au même nœud », par ex. Send("review", ReviewContext(doc=...)) vs Send("review", DefaultContext()).
  • Send équivalent à Command.goto : Command(goto=Send(...)) passe par le même canal TASKS dans map_command(goto as sends:62), donc conditional edges et valeur de retour de nœud convergent vers l'ordonnancement PUSH, abstraction unifiée.
  • Chaque Send est une task indépendante : N éléments dans une liste de Send deviennent N PregelExecutableTask indé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 — dataclass Send(node, arg, *, timeout=None), __slots__ = ("node", "arg", "timeout").
  • Send __init__:718 — constructeur, timeout passe par TimeoutPolicy.coerce.
  • Send __hash__:739hash((node, arg, timeout)), utilisé comme clé de déduplication dans le dict futures.
  • 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 canal TASKS, et appelle prepare_single_task pour assembler chaque PregelExecutableTask.
  • PULL tasks:470 — chemin PULL, calcule quels nœuds doivent s'exécuter selon trigger_to_nodes, ne consomme pas TASKS.
  • goto as sends:62map_command écrit les Send de Command.goto dans le canal TASKS, unifiant le chemin PUSH.
  • add_conditional_edges:969 — entrée des conditional edges, la fonction conditionnelle renvoie un Send ou une liste de Send.
  • PregelExecutableTask:627 — forme exécutable de la tâche, détient name / input / proc / writes / triggers / path etc.
  • PregelRunner.tick:176 — exécute en concurrence toutes les tâches PUSH + PULL de cette étape ; les tâches issues de Send y sont exécutées.

Flux de données

Send est réellement consommé au début de prepare_next_tasks(prepare_next_tasks:392) :

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)

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

  • Send doit cibler un nœud enregistré : si processes[send.node] n'existe pas, prepare_single_task renvoie None et 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é via add_node.
  • Send.arg ne passe pas par le reducer d'état : l'arg de Send sert 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 via config[CONFIG_KEY_READ] ou un paramètre state: Statearg n'est qu'une « entrée supplémentaire ».
  • Déduplication de liste de Send par hash : Send.__hash__ utilise le triplet (node, arg, timeout)(Send __hash__:739). Si deux Send vers le même nœud ont un arg non hachable (comme un dict), cela lève TypeErrorSend est 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 à foo alors qu'un canal trigger de foo est aussi mis à jour, prepare_next_tasks liste à la fois la tâche PUSH et la tâche PULL, et PregelRunner.tick les 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.timeout s'applique seule : Send(node, arg, timeout=10) ne s'applique qu'à cette tâche PUSH, sans affecter la configuration timeout par défaut du nœud cible(Send __init__:718).
  • Le canal TASKS est un Topic, pas un LastValue : Topic supporte l'accumulation de plusieurs écritures(Topic:23), donc plusieurs nœuds peuvent écrire des Send dans 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