Send: PUSH-Aufgaben-Primitiv zur Knoten-Ausfächerung
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
Sendexplizit übergeben).prepare_next_taskskonsumiert zuerst die PUSH-Aufgaben imTASKS-Kanal und berechnet dann die PULL-Aufgaben anhand vontrigger_to_nodes(PUSH tasks:440). argmuss nicht der Hauptzustand sein:Send.argkann ein beliebiges Objekt sein und muss nicht dem Hauptzustands-Schema des Graphen entsprechen — dasinput_schemades Zielknotens bestimmt die Form vonargund erlaubt «demselben Knoten heterogene Eingaben zu übergeben», z. B.Send("review", ReviewContext(doc=...))vsSend("review", DefaultContext()).SendundCommand.gotosind äquivalent:Command(goto=Send(...))läuft inmap_commandüber denselbenTASKS-Kanal(goto as sends:62); die beiden Einstiegspunkte (conditional edge und Knoten-Rückgabewert) münden also in derselben PUSH-Aufgaben-Scheduling-Abstraktion.- Jedes
Sendist eine eigenständige task:Die N Elemente einerSend-Liste sind N eigenständigePregelExecutableTaskmit 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— DatenklasseSend(node, arg, *, timeout=None),__slots__ = ("node", "arg", "timeout").Send __init__:718— Konstruktor;timeoutgeht überTimeoutPolicy.coerce.Send __hash__:739—hash((node, arg, timeout)), Grundlage für die Deduplizierung alsfutures-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 demTASKS-Kanal und ruft für jedesprepare_single_taskauf, um einePregelExecutableTaskzu bauen.PULL tasks:470— PULL-Pfad; berechnet anhand vontrigger_to_nodes, welche Knoten laufen sollen, und konsumiert denTASKS-Kanal nicht.goto as sends:62—map_commandschreibt dasSendausCommand.gotoin denTASKS-Kanal; vereinheitlicht den PUSH-Pfad.add_conditional_edges:969— Einstieg der conditional edge; die Bedingungsfunktion gibt einSendoder eineSend-Liste zurück.PregelExecutableTask:627— Ausführbare Form einer Aufgabe; hältname/input/proc/writes/triggers/pathusw.PregelRunner.tick:176— führt alle PUSH- + PULL-Aufgaben dieses Schritts gleichzeitig aus; hier werden die perSendübergebenen Aufgaben abgearbeitet.
Datenfluss
Send wird zu Beginn von prepare_next_tasks(prepare_next_tasks:392) konsumiert:
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
Sendmuss einen registrierten Zielknoten haben:Wennprocesses[send.node]nicht existiert, gibtprepare_single_taskNonezurü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 inadd_noderegistriert wurde.Send.arggeht nicht über den Zustands-Reducer:DasargvonSendist 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 überconfig[CONFIG_KEY_READ]oder den Parameterstate: Statedarauf zugreifen —argist nur «zusätzliche Eingabe».- Deduplizierung der
Send-Liste über hash:Send.__hash__verwendet das Tripel(node, arg, timeout)(Send __hash__:739); wenn diearg-Werte von zwei Sends zum selben Knoten nicht hashbar sind (wie dict), führt das direkt zuTypeError—Sendwird im Dedup-Pfad von_callals 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 vonfooaktualisiert wird, listetprepare_next_taskssowohl die PUSH- als auch die PULL-Aufgabe auf;PregelRunner.tickführt beide gleichzeitig aus — das ist erlaubt, aber Achtung: die Knotenfunktion kann im selben Schritt zweimal aufgerufen werden. Send.timeoutwirkt 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 istTopicnichtLastValue:Topicunterstützt kumulatives mehrfaches Schreiben(Topic:23), sodass im selben Superstep (superstep) mehrere KnotenSendin 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.