PregelRunner : exécuteur de tâches au sein d'un superpas
Responsabilités
PregelRunner ne fait qu'une seule chose : prendre le lot de PregelExecutableTask préparé par PregelLoop.tick et les exécuter réellement, puis collecter leurs writes et les rendre à PregelLoop. Il ne se soucie pas de la forme du graphe, ni du checkpoint, ni de la façon dont les canaux réduisent — il ne gère que « exécution concurrente + traitement des échecs + persistance des writes ».
Sa fenêtre est très étroite : dans la boucle principale de Pregel.astream (sync main loop:2964-2984), juste après que while loop.tick() a renvoyé True, vient for _ in runner.tick([t for t in loop.tasks.values() if not t.writes], ...):, suivi à la fin par loop.after_tick(). Autrement dit, l'entrée du runner, ce sont les « tâches calculées par tick qui n'ont pas encore écrit de writes » ; sa sortie, « les writes de ces tâches sont déjà entrés dans checkpoint_pending_writes » ; il n'y a pas de second tampon au milieu.
PregelRunner lui-même est très mince : deux méthodes publiques tick (tick:176) et commit (commit:574), plus une variante async atick. Le gros du travail est dans les deux appelants bas-niveau _call / _acall (_call:700)/(_acall:789) — ils enveloppent la fonction de nœud en future schedulable, gèrent la retry policy, le fan-out Send des sous-tâches, les callbacks de streaming chunk, et finissent par écrire les writes dans l'objet tâche.
Motivation de conception
Pourquoi le runner est-il un générateur plutôt qu'une méthode ordinaire ?
- Yield streaming : le runner est de la forme
for _ in runner.tick(...)(yield return type:188) ; chaque tâche terminée yield une fois, rendant la main à l'étage supérieur — qui peut alors pousser immédiatement le chunk streaming au client, sans attendre la fin du superpas. C'est la base du streaming LangGraph. - Chemin rapide pour tâche unique :
len(tasks) == 1 and timeout is None and get_waiter is Noneemprunte un appel synchrone direct, sans thread pool (fast path:203-254). La majorité des scénarios agents à branche unique tombent sur ce chemin et évitent le coût du pool. commitcomme callback de Future :FuturesDict.on_done(on_done:116) est déclenché à chaque future terminé ; en interne c'estweakref.WeakMethod(self.commit)—commitest le hook de persistance des writes par tâche, pas un nettoyage batch. Ainsi, l'ordre d'entrée des writes danspending_writessuit naturellement l'ordre d'achèvement, et le checkpoint peut s'écrire dès que possible.- Routage vers le handler d'erreur :
_should_route_to_error_handler(_should_route_to_error_handler:171-174) vérifie si le nœud en échec a un nœuderror_handlerconfiguré. Si oui, l'id de l'exception est ajouté à_handled_exception_idset une tâche de handler est planifiée à la place de la tâche originelle — l'échec ne panique pas immédiatement, il passe par un autre nœud pour le nettoyage. - Annulation via
_should_stop_others: dès qu'une tâche lève une exception nonGraphBubbleUp, le runner annule les autres tâches in-flight du même lot (_should_stop_others:616-634) — garantissant au sein d'un superpas la sémantique « tout réussit, ou on entre dans le handler ».GraphInterrupten est explicitement exclu, car une interruption n'est pas un échec.
Fichiers clés
class PregelRunner:135-138— définition de classe ; le docstring résume la responsabilité : exécuter les tâches, committer les writes, rendre la main, interrompre les autres tâches si besoin.__init__:140-169— détient des weakref verssubmit/put_writes, lenode_error_handler_map, les callbacksschedule_error_handler/aschedule_error_handler, et l'ensemble de déduplication d'exceptions cross-tick_handled_exception_ids.FuturesDict:75-134— sous-classe de dict personnalisée ; le callbackon_donedéclenchecommità chaque future terminé ; le champshould_stopporte le partial_should_stop_others.tick signature:176-188— entréeIterable[PregelExecutableTask], sortieIterator[None]; la signature expose la sémantique de générateur.fast path:203-254— quand tâche unique + pas de timeout + pas de waiter, exécution synchrone directe, avec planification du handler d'erreur en cas d'échec.schedule tasks:259-276— chemin multi-tâches : une future par tâche, schedulée viaself.submit()(en pratiquePregelLoop.submit).concurrent wait loop:282-323— boucleconcurrent.futures.wait(FIRST_COMPLETED); chaque tâche terminée déclenchecommit, émet les sorties, et au besoin lance une tâche handler.commit method:574-613— nettoyage par tâche : cancelled écrit ERROR,GraphInterruptécrit INTERRUPT, exception ordinaire écritERROR+ERROR_SOURCE_NODE, cas normal écrittask.writes+ marqueurNO_WRITES._should_stop_others:616-634— décide s'il faut annuler les autres tâches, en excluant explicitementGraphBubbleUpet les exceptions déjà handled._panic_or_proceed:650-697— fonction de nettoyage ; annule toutes les futures in-flight, fusionne plusieursGraphInterrupten un seul, lèveTimeoutErroren cas de timeout._call:700-787— enveloppe d'appel bas-niveau du nœud : retry, émission de stream chunk, callbackschedule_taskpour le fan-outSend, accumulation destask.writes._should_route_to_error_handler:171-174— renvoie True sitask.name in self.node_error_handler_map; le nœud handler d'erreur lui-même ne peut pas être routé (anti-récursion).
Flux de données
Le début de tick construit un FuturesDict, un dict amélioré qui appelle automatiquement commit à chaque future terminé :
def tick(
self,
tasks: Iterable[PregelExecutableTask],
*,
reraise: bool = True,
timeout: float | None = None,
retry_policy: Sequence[RetryPolicy] | None = None,
get_waiter: Callable[[], concurrent.futures.Future[None]] | None = None,
schedule_task: Callable[
[PregelExecutableTask, int, Call | None],
PregelExecutableTask | None,
],
) -> Iterator[None]:
tasks = tuple(tasks)
futures = FuturesDict(
callback=weakref.WeakMethod(self.commit),
event=threading.Event(),
should_stop=partial(
_should_stop_others, handled_exception_ids=self._handled_exception_ids
),
future_type=concurrent.futures.Future,
)
# give control back to the caller
yieldExtrait de tick 头部:176-199. Notez que le premier yield rend immédiatement la main — le premier for _ in runner.tick(...) de l'appelant n'est qu'un « démarrage » ; la véritable exécution des tâches a lieu dans les itérations suivantes. C'est une forme simplifiée de coroutine génératrice, pour éviter d'introduire le mot-clé async dans le chemin synchrone.
commit est la fonction de nettoyage par tâche, qui fait tomber les writes dans pending_writes :
def commit(
self,
task: PregelExecutableTask,
exception: BaseException | None,
) -> None:
if isinstance(exception, asyncio.CancelledError):
task.writes.append((ERROR, exception))
self.put_writes()(task.id, task.writes)
elif exception:
if isinstance(exception, GraphInterrupt):
if exception.args[0]:
writes = [(INTERRUPT, exception.args[0])]
if resumes := [w for w in task.writes if w[0] == RESUME]:
writes.extend(resumes)
self.put_writes()(task.id, writes)
elif isinstance(exception, GraphBubbleUp):
pass
else:
task.writes.append((ERROR, exception))
if self._should_route_to_error_handler(task) and not isinstance(
exception, GraphBubbleUp
):
task.writes.append((ERROR_SOURCE_NODE, task.name))
self._handled_exception_ids.add(id(exception))
self.put_writes()(task.id, task.writes)
else:
if self.node_finished and (
task.config is None or TAG_HIDDEN not in task.config.get("tags", [])
):
self.node_finished(task.name)
if not task.writes:
task.writes.append((NO_WRITES, None))
self.put_writes()(task.id, task.writes)Extrait de commit:574-613. put_writes est la weakref de PregelLoop.put_writes (put_writes:415) ; elle pousse les writes dans checkpoint_pending_writes et lance la future de persistance du checkpointer. Notez les writes spéciaux :
CancelledError→ écrit(ERROR, exception), pour qu'à la fermeture d'after_tickla tâche soit comptée comme « échec mais enregistré », sans panic.GraphInterrupt→ écrit(INTERRUPT, ...)+ tous les writesRESUME; la valeur de reprise de l'interruption doit être persistée au côté de l'interruption elle-même, pour être relue ensemble au prochain resume.- Exception ordinaire + handler d'erreur configuré → écrit en plus
(ERROR_SOURCE_NODE, task.name); ce marqueur sera scanné parPregelLoop._resume_error_handlers_if_applicable(ERROR_SOURCE_NODE scan:775) pour planifier le nœud handler correspondant à la reprise. - Terminaison normale sans rien écrire → ajoute
(NO_WRITES, None), pour empêcherprepare_next_tasksde la prendre pour « pas encore exécutée » (NO_WRITES marker:609-611).
Le flux d'exécution complet du runner se dessine ainsi :
Limites et échecs
_handled_exception_idspersiste entre ticks : dans__init__, cet ensemble est au niveau instance, pas par tick (_handled_exception_ids:169). Un même id d'objet exception entré dans l'ensemble ne sera plus considéré comme échec par_should_stop_others, évitant que le handler d'erreur finisse par relancer l'exception originelle.GraphBubbleUpn'est pas un échec :commitfait explicitementpasssurGraphBubbleUp, sans écrire aucun write (GraphBubbleUp pass:592-594). Ce type d'exception (par ex.ParentCommand) sert à envoyer un signal au graphe parent, ce n'est pas une vraie erreur.NO_WRITESévite la ré-exécution de tâche : une tâche qui termine normalement sans rien écrire doit quand même recevoir un(NO_WRITES, None)(NO_WRITES:609-611). Sinon, aprèsafter_tick,prepare_next_tasksverrait cette tâche sans writes, la croirait non exécutée, et la relancerait à la reprise.schedule_error_handlerpeut renvoyer None : le nœud handler lui-même peut ne pas être configuré ou impossible à construire ; la tâche emprunte alors le chemin panic ordinaire (handler optional:230-248). Le writeERROR_SOURCE_NODEdecommitcoexiste avec ce chemin, et_resume_error_handlers_if_applicableretentera la planification à la reprise.- Un timeout ne lève pas GraphInterrupt : dans
_panic_or_proceed,if inflight: raise timeout_exc_cls("Timed out")(timeout:691-697) — timeout est une vraie exception, elle remonte hors de Pregel, contrairement à GraphInterrupt qui est avalé. - Élagage du traceback : avant le reraise, on saute les frames listées dans
EXCLUDED_FRAME_FNAMES(tb trim:241-247) pour filtrer les frames internes à langgraph, afin que le traceback vu par l'utilisateur pointe directement vers son code de nœud.
Résumé
PregelRunner est une coque mince : le générateur tick planifie l'exécution concurrente + le commit par tâche fait tomber les writes dans pending_writes, avec au milieu le routage vers le handler d'erreur et le mécanisme d'annulation. Une fois compris, on a compris « comment une tâche est réellement exécutée au sein d'un superpas ». Pour le pilote de boucle PregelLoop à l'étage supérieur, voir /pregel/loop ; pour les trois fonctions algorithmiques prepare_next_tasks / apply_writes / should_interrupt, voir /pregel/algo ; pour l'assemblage global, voir /pregel/pregel ; pour la façon dont le streaming est façonné entre les yields du runner, voir /stream/run-stream.
Voir la documentation officielle : LangGraph docs · README.