PregelLoop : pilote de la boucle de superpas
Responsabilités
PregelLoop est le « métronome » du runtime LangGraph. Pregel lui-même ne fait qu'assembler les ressources et exposer les entrées invoke / astream ; ce qui fait tourner le graphe, cercle après cercle, est une instance de PregelLoop : elle détient le point de contrôle (checkpoint) courant, l'ensemble des canaux (channels), les pending_writes, le compteur step, la machine à états status, ainsi qu'un dictionnaire tasks — tout le contexte d'un seul superpas (superstep). Le schéma enveloppant while loop.tick(): ... loop.after_tick() s'appuie sur elle.
Ses deux méthodes centrales sont tick (tick:599) et after_tick (after_tick:683). La première gère le « avant exécution » : vérifier si l'on a dépassé la limite de pas, appeler prepare_next_tasks pour charger dans self.tasks les nœuds à exécuter pour ce pas, examiner les interruptions, traiter drain_requested, et recoller les pending_writes d'une reprise précédente sur les nœuds déjà réussis. La seconde gère le « après exécution » : collecter tous les task.writes, appeler apply_writes pour les fusionner dans les canaux, nettoyer les pending, poser le checkpoint, puis décider des interrupt_after. L'une et l'autre encadrent le cycle BSP « lire l'entrée → planifier → exécuter → réécrire → persister ».
Autrement dit, PregelLoop n'exécute pas elle-même les nœuds — c'est le travail de PregelRunner ; elle décide seulement « qui regarder à ce pas, comment finaliser après exécution, et s'il faut continuer au pas suivant ».
Motivation de conception
Pourquoi séparer la logique de boucle en tick et after_tick plutôt que d'embrayer tout dans une méthode step() ?
- Rendre l'exécution à l'appelant : une fois que
tickrenvoieTrue, l'appelant (le corps dePregel.astream, voirsync main loop:2964-2984) pilote lui-mêmerunner.tickpour faire tourner ce batch de tâches, puis rappelleloop.after_tick.PregelLoopn'a ainsi pas besoin de détenir une référence àPregelRunner; les deux objets collaborent en couplage lâche, sans imbrication. - Machine à états explicite : le champ
statusprend sept valeurs —"input" / "pending" / "done" / "draining" / "interrupt_before" / "interrupt_after" / "out_of_steps"(status Literal:256-264) ; chaque état correspond à une raison de sortie, le code externe s'en sert pour décider s'il faut resume, raise ou break, sans devoir deviner la sémantique de la valeur de retour. - Sémantique claire de reprise après rupture :
is_replayingn'est vrai qu'au premier tick — il indique que ce pas « rejoue » les tâches déjà achevées avant la dernière interruption ; dès queafter_ticktermine, il est désactivé (is_replaying = False:716). Combiné à_reapply_writes_to_succeeded_nodesqui recolle les writes réussis du checkpoint sur les tâches mémoire, cela permet aux nœuds réussis de ne pas être relancés et aux nœuds en échec d'entrer dans le handler d'erreur. - Drain est la porte d'intervention externe :
RunControl.drain_requestedest positionné par le consommateur du stream (par ex. un client qui veut s'arrêter en cours) ; dès quetickle voit, il passe à l'étatdraininget renvoieFalse, lewhileexterne sort naturellement, sans avoir à tuer un nœud en plein milieu. - Écriture en mode delta et timing du checkpoint :
_delta_write_futs(_delta_write_futs:207) collecte les futuresput_writesde tous les canaux delta ;_checkpointer_put_after_previousdoit d'abord les drainer avant d'écrire le checkpoint suivant — pour garantir l'ordre causal « l'écriture produit le checkpoint, pas le checkpoint qui écrase l'écriture ».
Fichiers clés
class PregelLoop:158— définition de classe ; les champs sont quasiment tous des attributs bruts, sans property, pour que les sous-classesSyncPregelLoop/AsyncPregelLooppuissent les surcharger directement.status Literal:256-264— énumération à sept états ;"out_of_steps"est déclenché par la limite de pas,"draining"par une requête d'arrêt externe,"done"est la fin normale.tick method:599-681— toute la logique d'un tick complet : vérification du pas,prepare_next_tasks, branche done sur tâche vide, drain,_reapply_writes_to_succeeded_nodes,should_interrupt.out_of_steps:607-609— quandself.step > self.stop, passe à l'étatout_of_stepset renvoieFalse, lewhileexterne sort.done branch:653-655— quandself.tasksest vide, passe à l'étatdone, c'est la « mort naturelle » du graphe.draining branch:657-659— quandcontrol.drain_requestedest vrai, passe à l'étatdraining.reapply + resume handlers:662-664— à la reprise, on colle d'abord les pending writes sur les nœuds réussis, puis on appelle_resume_error_handlers_if_applicablepour préparer des tâches de handler pour les nœuds en échec.interrupt_before:667-671— appelleshould_interruptpour vérifier si un nœud tombe dans la listeinterrupt_before; si oui, lèveGraphInterrupt.after_tick method:683-726— collecte les writes →apply_writes→ émet values → nettoie pending →_put_checkpoint→ vérificationinterrupt_after._reapply_writes_to_succeeded_nodes:736-749— recolle les pending writes sur les tâches mémoire, mais saute les quatre signaux de contrôleERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME._resume_error_handlers_if_applicable:751-816— scanne le marqueurERROR_SOURCE_NODE, construit une tâche de handler d'erreur pour le nœud en échec, afin que le runner saute la tâche originelle et exécute directement le handler.put_writes:415-508— entrée des writes produites par une tâche ; déduplique les canaux de contrôle, accumule les null task, lance les futuresput_writesdu checkpointer.SyncPregelLoop:1469/AsyncPregelLoop:1722— deux sous-classes, surchargent__enter__/__exit__,accept_push,schedule_error_handleret les autres parties différant entre sync et async.
Flux de données
L'entrée de tick effectue d'abord la vérification de la limite de pas, puis appelle prepare_next_tasks pour calculer « quelles tâches exécuter au pas suivant ». C'est la graine de tout le superpas — toutes les opérations ultérieures travaillent ce batch de tâches.
def tick(self) -> bool:
"""Execute a single iteration of the Pregel loop.
Returns:
True if more iterations are needed.
"""
# check if iteration limit is reached
if self.step > self.stop:
self.status = "out_of_steps"
return False
# prepare next tasks
self.tasks = prepare_next_tasks(
self.checkpoint,
self.checkpoint_pending_writes,
self.nodes,
self.channels,
self.managed,
self.config,
self.step,
self.stop,
for_execution=True,
manager=self.manager,
store=self.store,
checkpointer=self.checkpointer,
trigger_to_nodes=self.trigger_to_nodes,
updated_channels=self.updated_channels,
retry_policy=self.retry_policy,
cache_policy=self.cache_policy,
)Extrait de tick 头部:599-629. self.step > self.stop est la condition de sortie « limite de pas » — stop est fixé à recursion_limit dans PregelLoop.__init__ ; si la profondeur de récursion est dépassée, on s'arrête sur out_of_steps. prepare_next_tasks prend le checkpoint courant, les pending_writes, la table des nœuds et la table des canaux pour calculer quelles tâches exécuter à ce pas et les ranger dans self.tasks. for_execution=True indique que ces tâches doivent réellement s'exécuter (et non un dry-run pour aperçu streaming).
Viennent ensuite trois conditions de sortie en parallèle — tâche vide / drain / reprise collant les writes — puis should_interrupt :
# if no more tasks, we're done
if not self.tasks:
self.status = "done"
return False
if self.control is not None and self.control.drain_requested:
self.status = "draining"
return False
# if there are pending writes from a previous loop, apply them
if not self.is_replaying and self.checkpoint_pending_writes:
self._reapply_writes_to_succeeded_nodes(self.tasks)
self._resume_error_handlers_if_applicable()
# before execution, check if we should interrupt
if self.interrupt_before and should_interrupt(
self.checkpoint, self.interrupt_before, self.tasks.values()
):
self.status = "interrupt_before"
raise GraphInterrupt()Extrait de tick 中段:652-671. La vérification interrupt_before n'a lieu qu'après re-collage des pending_writes en cas de reprise — car à la reprise les versions de canaux reflètent déjà les écritures précédentes, ce qui permet à should_interrupt de décider correctement « y a-t-il eu de nouvelles mises à jour depuis la dernière interruption ».
after_tick est appelé quand le runner a terminé toutes les tâches ; il gère le nettoyage :
def after_tick(self) -> None:
# finish superstep
writes = [w for t in self.tasks.values() for w in t.writes]
self._delta_channels_with_overwrite.update(
ch
for ch, v in writes
if isinstance(self.specs.get(ch), DeltaChannel) and _get_overwrite(v)[0]
)
# all tasks have finished
self.updated_channels = apply_writes(
self.checkpoint,
self.channels,
self.tasks.values(),
self.checkpointer_get_next_version,
self.trigger_to_nodes,
)Extrait de after_tick 头部:683-698. On aplatit d'abord les writes de toutes les tâches, en notant au passage quels canaux delta ont subi un overwrite (ce qui affecte la valeur de départ du replay sparse) ; puis on appelle apply_writes pour fusionner ces writes dans les canaux, et la valeur renvoyée updated_channels servira de hint au prepare_next_tasks suivant pour accélérer « trouver le prochain nœud à déclencher ». La fin de after_tick vide aussi checkpoint_pending_writes, passe is_replaying à False, pose le checkpoint via _put_checkpoint({"source": "loop"}), puis exécute le test interrupt_after (after_tick 尾部:714-724).
Le cycle de vie complet d'un superpas se dessine ainsi :
Limites et échecs
out_of_stepsn'est pas une erreur :step > stoprenvoie directementFalse, le status est mis àout_of_steps, lewhileexterme sort naturellement. L'étage supérieur dePregeldécide selon le status s'il faut leverRecursionErrorou s'arrêter silencieusement — c'est une sémantique de « plafond souple », pas un chemin d'exception (out_of_steps:607-609).drainingne doit pas émettre d'événements lifecycle :_emit_graph_lifecycle_eventrefuse explicitement d'être appelé à l'étatdraining(draining guard:384-385), car drain est un arrêt volontaire de l'utilisateur, et n'a pas la sémantique d'interrupt / resume.- À la reprise, les pending writes contenant des signaux de contrôle doivent être sautés :
_reapply_writes_to_succeeded_nodesdoit ignorerERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME(skip control signals:746-747), sinon les tâches qui ont échoué seraient considérées comme réussies, et les handlers déjà joués seraient redéclenchés. - Re-planification du handler d'erreur au resume :
_resume_error_handlers_if_applicablene s'applique qu'aux nœuds disposant d'unerror_handler_node(handler_node check:791-793) ; les nœuds en échec sans handler restent en l'état, et le_should_stop_othersdu runner les emmène sur le chemin panic. - Sémantique d'accumulation du null task dans
put_writes: les writes deNULL_TASK_IDne s'écrasent pas mais s'accumulent (null task accumulate:422-431), pour les écritures de type input qui « n'appartiennent à aucun nœud mais doivent entrer dans le checkpoint ». - Les futures du canal delta doivent être flushés avant le prochain checkpoint :
_delta_write_futscollecte les futuresput_writesde tous les canaux delta ; avant d'écrire le checkpoint suivant, on les draine (_delta_write_futs:201-207) — inverser l'ordre ferait que le replay sparse ne verrait pas les writes qui ont produit le checkpoint.
Résumé
PregelLoop est un métronome de superpas piloté par une machine à états, qui encadre le triptyque « préparer les tâches → exécuter → nettoyer » autour de PregelRunner sans toucher lui-même à l'exécution. Comprendre les deux segments tick / after_tick, les sept états de status, et le mécanisme de reprise is_replaying + _reapply_writes_to_succeeded_nodes, c'est saisir l'axe principal du modèle d'exécution de LangGraph. Pour la façon dont le runner fait réellement tourner les tâches, voir /pregel/runner ; pour les détails internes des fonctions algorithmiques prepare_next_tasks / apply_writes / should_interrupt, voir /pregel/algo ; pour l'assemblage global, voir /pregel/pregel. La sémantique d'écriture des canaux est dans /channel/base-channel.
Voir la documentation officielle : LangGraph docs · README.