Skip to content

PregelRunner : exécuteur de tâches au sein d'un superpas

源码版本1.2.9

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 None emprunte 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.
  • commit comme callback de Future : FuturesDict.on_done (on_done:116) est déclenché à chaque future terminé ; en interne c'est weakref.WeakMethod(self.commit)commit est le hook de persistance des writes par tâche, pas un nettoyage batch. Ainsi, l'ordre d'entrée des writes dans pending_writes suit 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œud error_handler configuré. Si oui, l'id de l'exception est ajouté à _handled_exception_ids et 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 non GraphBubbleUp, 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 ». GraphInterrupt en 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 vers submit / put_writes, le node_error_handler_map, les callbacks schedule_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 callback on_done déclenche commit à chaque future terminé ; le champ should_stop porte le partial _should_stop_others.
  • tick signature:176-188 — entrée Iterable[PregelExecutableTask], sortie Iterator[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 via self.submit() (en pratique PregelLoop.submit).
  • concurrent wait loop:282-323 — boucle concurrent.futures.wait(FIRST_COMPLETED) ; chaque tâche terminée déclenche commit, é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 écrit ERROR + ERROR_SOURCE_NODE, cas normal écrit task.writes + marqueur NO_WRITES.
  • _should_stop_others:616-634 — décide s'il faut annuler les autres tâches, en excluant explicitement GraphBubbleUp et les exceptions déjà handled.
  • _panic_or_proceed:650-697 — fonction de nettoyage ; annule toutes les futures in-flight, fusionne plusieurs GraphInterrupt en un seul, lève TimeoutError en cas de timeout.
  • _call:700-787 — enveloppe d'appel bas-niveau du nœud : retry, émission de stream chunk, callback schedule_task pour le fan-out Send, accumulation des task.writes.
  • _should_route_to_error_handler:171-174 — renvoie True si task.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é :

python
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
    yield

Extrait 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 :

python
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_tick la tâche soit comptée comme « échec mais enregistré », sans panic.
  • GraphInterrupt → écrit (INTERRUPT, ...) + tous les writes RESUME ; 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é par PregelLoop._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êcher prepare_next_tasks de 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_ids persiste 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.
  • GraphBubbleUp n'est pas un échec : commit fait explicitement pass sur GraphBubbleUp, 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ès after_tick, prepare_next_tasks verrait cette tâche sans writes, la croirait non exécutée, et la relancerait à la reprise.
  • schedule_error_handler peut 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 write ERROR_SOURCE_NODE de commit coexiste avec ce chemin, et _resume_error_handlers_if_applicable retentera 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.