Skip to content

Send: PUSH-Aufgaben-Primitiv zur Knoten-Ausfächerung

源码版本1.2.9

Verantwortung

Send ist LangGraphs Primitiv für «Aufgabe an einen bestimmten Knoten übergeben», definiert in libs/langgraph/langgraph/types.py(Send class:664). Ein Send(node="foo", arg={"k": v}) bedeutet «im nächsten Schritt den Knoten foo mit arg als Eingabe aufrufen». Der Unterschied zur normalen Knoten-Auslösung: Ein normaler Knoten (PULL) wird passiv durch Kanal-Abonnements ausgelöst, ein Send (PUSH) ist eine aktive Erklärung des Vorgängerknotens: «Ich möchte, dass foo einmal läuft, mit diesem spezifischen arg».

Seine Position liegt zwischen conditional edge und dem TASKS-Kanal:add_conditional_edges(START, lambda s: [Send("gen", {"x": 1}), Send("gen", {"x": 2})]) schreibt die von der Bedingungsfunktion zurückgegebene Send-Liste in den TASKS-Kanal(goto as sends:62); die nächste Runde prepare_next_tasks holt sie in tasks_channel.get() heraus und baut jedes Send zu einer eigenständigen PregelExecutableTask zusammen(PUSH tasks:421).

Entwurfsmotivation

  • Natürliche Darstellung von Map-Reduce:Das klassische Szenario ist «denselben Knoten parallel N-mal laufen zu lassen, jedes Mal mit anderer Eingabe». Send("generate", {"subject": s}) gibt in einer conditional edge eine Liste zurück, jedes Send ist ein eigenständiger Aufruf, Pregel führt sie von Natur aus gleichzeitig aus, und der Reducer fasst die N Ergebnisse in den Hauptzustand zusammen. Ein manuelles Ausfächern von N Kanten im Graphen entfällt.
  • PUSH / PULL-Dichotomie:Die Auslösung eines Knotens zerfällt in zwei Klassen — PULL (passiv, durch Updates der abonnierten Kanäle ausgelöst) und PUSH (aktiv, durch Send explizit übergeben). prepare_next_tasks konsumiert zuerst die PUSH-Aufgaben im TASKS-Kanal und berechnet dann die PULL-Aufgaben anhand von trigger_to_nodes(PUSH tasks:440).
  • arg muss nicht der Hauptzustand sein:Send.arg kann ein beliebiges Objekt sein und muss nicht dem Hauptzustands-Schema des Graphen entsprechen — das input_schema des Zielknotens bestimmt die Form von arg und erlaubt «demselben Knoten heterogene Eingaben zu übergeben», z. B. Send("review", ReviewContext(doc=...)) vs Send("review", DefaultContext()).
  • Send und Command.goto sind äquivalent:Command(goto=Send(...)) läuft in map_command über denselben TASKS-Kanal(goto as sends:62); die beiden Einstiegspunkte (conditional edge und Knoten-Rückgabewert) münden also in derselben PUSH-Aufgaben-Scheduling-Abstraktion.
  • Jedes Send ist eine eigenständige task:Die N Elemente einer Send-Liste sind N eigenständige PregelExecutableTask mit eigenem task_id und path, die unabhängig retried und fehlerbehandelt werden können — ein Ausfall eines Send beeinträchtigt nicht die anderen gleichzeitigen Sends.

Schlüsseldateien

  • Send class:664 — Datenklasse Send(node, arg, *, timeout=None), __slots__ = ("node", "arg", "timeout").
  • Send __init__:718 — Konstruktor; timeout geht über TimeoutPolicy.coerce.
  • Send __hash__:739hash((node, arg, timeout)), Grundlage für die Deduplizierung als futures-dict-key.
  • prepare_next_tasks:392 — Haupt-Scheduling-Funktion; konsumiert zuerst PUSH-Aufgaben und berechnet dann die PULL-Kandidaten.
  • PUSH tasks:440 — holt die Send-Sequenz aus dem TASKS-Kanal und ruft für jedes prepare_single_task auf, um eine PregelExecutableTask zu bauen.
  • PULL tasks:470 — PULL-Pfad; berechnet anhand von trigger_to_nodes, welche Knoten laufen sollen, und konsumiert den TASKS-Kanal nicht.
  • goto as sends:62map_command schreibt das Send aus Command.goto in den TASKS-Kanal; vereinheitlicht den PUSH-Pfad.
  • add_conditional_edges:969 — Einstieg der conditional edge; die Bedingungsfunktion gibt ein Send oder eine Send-Liste zurück.
  • PregelExecutableTask:627 — Ausführbare Form einer Aufgabe; hält name / input / proc / writes / triggers / path usw.
  • PregelRunner.tick:176 — führt alle PUSH- + PULL-Aufgaben dieses Schritts gleichzeitig aus; hier werden die per Send übergebenen Aufgaben abgearbeitet.

Datenfluss

Send wird zu Beginn von prepare_next_tasks(prepare_next_tasks:392) konsumiert:

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)

Der TASKS-Kanal ist ein Topic[Send] — ein Kanaltyp, der mehrfaches Schreiben und Lesen unterstützt(Topic:23). Jedes Send wird aus tasks_channel.get() geholt; (PUSH, idx) dient als task_id-Präfix für prepare_single_task, das aus processes[send.node] die entsprechende PregelNode-Konfiguration (inkl. proc / writers / retry_policy usw.) holt und mit send.arg als input eine PregelExecutableTask zusammenbaut. Send.arg geht nicht in den Hauptzustandskanal — es wird direkt als task.input an den proc des Knotens übergeben; die Knotenfunktion erhält dieses arg und nicht den vollständigen Snapshot des Hauptzustands des Graphen.

Es gibt zwei Pfade, auf denen Send in den TASKS-Kanal gelangt: (1) Die in add_conditional_edges konfigurierte Bedingungsfunktion gibt eine Send-Liste zurück, und attach_branch schreibt sie in TASKS; (2) die Knotenfunktion gibt Command(goto=Send(...)) zurück, und map_command in _io.py schreibt das Send in TASKS(goto as sends:62). Beide Pfade münden im selben TASKS-Kanal und werden in der nächsten Runde von prepare_next_tasks konsumiert.

Grenzen und Fehler

  • Send muss einen registrierten Zielknoten haben:Wenn processes[send.node] nicht existiert, gibt prepare_single_task None zurück und das Send wird still verworfen — es wird nicht raisen, aber die Aufgabe läuft nicht. Wenn beim Debugging auffällt, dass ein Send nicht wirkt, zuerst prüfen, ob der Zielknotenname in add_node registriert wurde.
  • Send.arg geht nicht über den Zustands-Reducer:Das arg von Send ist direkt der input des Knotens und wird nicht über den Reducer des Hauptzustandskanals zusammengeführt. Wenn der Zielknoten den Hauptzustand lesen möchte, muss er in der Knotenfunktion über config[CONFIG_KEY_READ] oder den Parameter state: State darauf zugreifen — arg ist nur «zusätzliche Eingabe».
  • Deduplizierung der Send-Liste über hash:Send.__hash__ verwendet das Tripel (node, arg, timeout)(Send __hash__:739); wenn die arg-Werte von zwei Sends zum selben Knoten nicht hashbar sind (wie dict), führt das direkt zu TypeErrorSend wird im Dedup-Pfad von _call als dict-key verwendet.
  • Derselbe Knoten kann gleichzeitig per PULL und PUSH ausgelöst werden:Wenn Send("foo", arg) foo übermittelt und gleichzeitig ein trigger-Kanal von foo aktualisiert wird, listet prepare_next_tasks sowohl die PUSH- als auch die PULL-Aufgabe auf; PregelRunner.tick führt beide gleichzeitig aus — das ist erlaubt, aber Achtung: die Knotenfunktion kann im selben Schritt zweimal aufgerufen werden.
  • Send.timeout wirkt isoliert:Send(node, arg, timeout=10) gilt nur für diese PUSH-Aufgabe und beeinflusst nicht die Standard-timeout-Konfiguration des Zielknotens(Send __init__:718).
  • TASKS-Kanal ist Topic nicht LastValue:Topic unterstützt kumulatives mehrfaches Schreiben(Topic:23), sodass im selben Superstep (superstep) mehrere Knoten Send in den Kanal schreiben können und die nächste Runde alles auf einmal konsumiert — das ist die Grundlage von Map-Reduce.

Zusammenfassung

Send ist der Träger von PUSH-Aufgaben und ergänzt die normalen PULL-Knoten. Die beiden Einstiegspunkte (conditional edge und Command.goto) münden in demselben TASKS-Kanal und werden einheitlich von prepare_next_tasks schedult. Wie tasks parallel ausgeführt werden, siehe /func/concurrency; die Beziehung zwischen conditional edge und Knoten siehe libs/langgraph/langgraph/graph/_branch.py. Siehe offizielle Dokumentation: LangGraph 文档 · README.