Couche algorithmique de Pregel : interruption, écriture, planification
Responsabilités
_algo.py est le « noyau mathématique » de PregelLoop et PregelRunner. La loop s'occupe de « combien de tours », le runner de « comment exécuter en parallèle », et l'algo s'occupe de « la sémantique de chaque pas » : quand interrompre, comment fusionner les writes d'un nœud dans les canaux (channels), et quels nœuds déclencher au pas suivant. Pas d'état, pas d'IO — que des fonctions pures ; en entrée checkpoint + canaux + tâches, en sortie nouvelle table de versions de canaux / ensemble updated_channels / un nouveau dictionnaire de tâches.
Les trois fonctions centrales correspondent exactement aux trois appels de PregelLoop.tick + after_tick : should_interrupt (should_interrupt:155) appelé au milieu de tick, décide s'il faut lever GraphInterrupt avant exécution ; apply_writes (apply_writes:232) appelé au début d'after_tick, fusionne dans les canaux les writes de toutes les tâches du pas et renvoie updated_channels ; prepare_next_tasks (prepare_next_tasks:392) appelé au début de tick, calcule à partir du checkpoint courant quelles tâches exécuter à ce pas.
Autrement dit, ces trois fonctions décident « ce que Pregel voit, ce qu'il fait, et l'état des canaux après nettoyage, à chaque pas de la boucle BSP ». Les comprendre, c'est comprendre les « règles de transition d'état » de Pregel.
Motivation de conception
Pourquoi découper la logique algorithmique en fonctions pures ?
- Testabilité :
should_interrupt/apply_writes/prepare_next_taskssont toutes pures ; en entrée un dictionnaire checkpoint + un mapping de canaux, en sortie une structure de données. Elles peuvent être testées unitairement sans l'IO, le submit, le checkpointer de Pregel — c'est précisément la forme de la plupart des tests delibs/langgraph/tests/unit/. - Partage sync/async :
PregelLoopa deux sous-classesSyncPregelLoopetAsyncPregelLoop, mais la couche algo est sans side effect, les deux sous-classes partagent la même logique (tick calls prepare_next_tasks:612). Sinon, il faudrait écrire l'algorithme en deux exemplaires sync/async, et corriger les bug des deux côtés. - Optimisation
updated_channels:apply_writesrenvoie l'ensemble updated_channels, etprepare_next_tasksle reçoit comme hint, sautant « scanner tous les nœuds pour voir si trigger est déclenché » (updated_channels hint:475-486). Sur un grand graphe, cette étape fait passer le scan de trigger en O(N×M) à O(updated×triggered), c'est l'optimisation de perf la plus importante. should_interruptdéduplique viaversions_seen: leversions_seendu canalINTERRUPT(versions_seen INTERRUPT:163) enregistre le numéro de version des canaux au moment de la dernière interruption. On ne considère une nouvelle interruption que si « des mises à jour de canaux sont survenues depuis la dernière interruption » — pour éviter la boucle infinie interrupt → resume → interrupt.- Signal de fin à la fin d'
apply_writes: quandbump_step and updated_channels.isdisjoint(trigger_to_nodes), on appellechannels[chan].finish()(finish:336-342) pour envoyer à tous les canaux le signal « c'est le dernier superpas ». Les canaux permanents (commeLastValue) en profitent pour se marquer indisponibles, pour queprepare_next_tasksne déclenche plus aucun nœud — c'est la « mort naturelle » du graphe.
Fichiers clés
should_interrupt:155-185— vérifie si le graphe doit interrompre ; regarde d'abordversions_seen[INTERRUPT]pour savoir s'il y a eu de nouvelles mises à jour, puis si la tâche déclenchée tombe dans la listeinterrupt_nodes.any_updates_since_prev_interrupt:161-168— compare les numéros de version des canaux pour savoir « y a-t-il eu une mise à jour depuis la dernière interruption » ; c'est le cœur de la déduplication des interruptions.interrupt_nodes filter:170-184—interrupt_nodes == "*"signifie que toutes les tâches non hidden sont prises en compte ; sinon on filtre partask.name in interrupt_nodes.apply_writes:232-345— fusionne dans les canaux les writes des tâches du pas, renvoie l'ensembleupdated_channels.bump_step:256-259— si une tâche a destriggers,bump_step=True, ce pas consomme des canaux et fait avancer le compteur de step.update seen versions:262-269— chaque tâche enregistre dansversions_seenle numéro de version des canaux trigger qu'elle a lus ; au prochain_triggers, on comparera pour savoir « y a-t-il du neuf depuis la dernière exécution ».consume channels:284-292— appellechannels[chan].consume()pour marquer comme consommés les canaux lus à ce pas ; c'est l'implémentation de la sémantique « lecture unique efface » des canaux PULL.apply writes to channels:315-323—channels[chan].update(vals)fusionne réellement les writes dans les canaux ; seuls les canauxis_available()entrent dansupdated_channels.bump_step notify:326-333— quandbump_step=True, les canaux disponibles non mis à jour à ce pas reçoivent aussiupdate(EMPTY_SEQ), pour faire avancer leur numéro de version — ainsi les nœuds abonnés à ces canaux mais non déclenchés à ce pas ne seront pas fautivement réveillés au prochain tick.finish signal:336-342— si aucun nœud n'a été déclenché à ce pas, envoiefinish()à tous les canaux ; les canaux permanents en profitent pour devenir unavailable, et le graphe se termine naturellement.prepare_next_tasks:392-513— calcule quelles tâches exécuter au pas suivant ; PUSH (issu du fan-outSend) et PULL (issu des déclencheurs d'arêtes) en sortent tous deux.updated_channels optimization:475-486— utilise leupdated_channelsdu pas précédent +trigger_to_nodespour inverser la recherche des nœuds déclenchés, en sautant le scan complet.prepare_single_task:524— construction d'une tâche unique ; PUSH passe parprepare_push_task_functional/prepare_push_task_send, PULL par le test_triggers+_proc_inputpour l'entrée._triggers:1260-1277— teste si les canaux trigger du nœud ont « une version non lue » ; c'est la condition centrale de la planification PULL.local_read:188-224— entrée de la lecture des canaux par la fonction de nœud ; supportefresh=Truepour lire la vue locale « ce que je viens d'écrire mais pas encore fusionné dans les canaux », utilisée pour la sémantique « écrit-then-lit immédiat » de l'arête conditionnelle (conditional edge).
Flux de données
should_interrupt est la plus simple des trois ; sa logique est « il y a du neuf + la tâche courante est dans la liste » :
def should_interrupt(
checkpoint: Checkpoint,
interrupt_nodes: All | Sequence[str],
tasks: Iterable[PregelExecutableTask],
) -> list[PregelExecutableTask]:
"""Check if the graph should be interrupted based on current state."""
version_type = type(next(iter(checkpoint["channel_versions"].values()), None))
null_version = version_type() # type: ignore[misc]
seen = checkpoint["versions_seen"].get(INTERRUPT, {})
# interrupt if any channel has been updated since last interrupt
any_updates_since_prev_interrupt = any(
version > seen.get(chan, null_version) # type: ignore[operator]
for chan, version in checkpoint["channel_versions"].items()
)
# and any triggered node is in interrupt_nodes list
return (
[task for task in tasks if ...]
if any_updates_since_prev_interrupt
else []
)Extrait de should_interrupt:155-185. La clé est versions_seen[INTERRUPT] — ce n'est pas le numéro de version vu par un nœud, c'est « un instantané de tous les numéros de version de canaux au moment de la dernière interruption ». any_updates_since_prev_interrupt compare cet instantané au channel_versions courant ; si un canal a une version plus haute, on estime qu'« il y a du neuf », et seulement alors on entre dans le test d'interruption. Ainsi, après un resume, même si le tick déclenche des nœuds du même nom, il n'y aura pas immédiatement une nouvelle interruption — il faut attendre que les canaux aient réellement reçu de nouvelles écritures.
apply_writes est le cœur du nettoyage ; on met d'abord à jour versions_seen, puis on consomme les canaux lus, et enfin on écrit les nouvelles valeurs :
# update seen versions
for task in tasks:
checkpoint["versions_seen"].setdefault(task.name, {}).update(
{
chan: checkpoint["channel_versions"][chan]
for chan in task.triggers
if chan in checkpoint["channel_versions"]
}
)
# Consume all channels that were read
for chan in {
chan
for task in tasks
for chan in task.triggers
if chan not in RESERVED and chan in channels
}:
if channels[chan].consume() and next_version is not None:
checkpoint["channel_versions"][chan] = next_version
# Apply writes to channels
updated_channels: set[str] = set()
for chan, vals in pending_writes_by_channel.items():
if chan in channels:
if channels[chan].update(vals) and next_version is not None:
checkpoint["channel_versions"][chan] = next_version
if channels[chan].is_available():
updated_channels.add(chan)Extrait de apply_writes 核心:262-323. La sémantique en trois segments est limpide :
- Mise à jour de
versions_seen— on enregistre le numéro de version courant des canaux trigger de chaque tâche, comme preuve que « ce nœud a vu cette version de canal ». Au prochain test_triggers, on comparera avec les versions futures, et seul un delta déclenchera un nouveau run. consume()les canaux lus — les canaux « lecture unique efface » commeTopicvident leur buffer interne dansconsume(), en ne conservant que l'avance du numéro de version. C'est le mécanisme par lequel les messages fan-out parSend, une fois consommés, ne redéclenchent pas d'autres nœuds.update(vals)écrit les nouvelles valeurs — les writes de toutes les tâches du pas pour chaque canal sont collectés et fusionnés en une fois, déclenchant la logique de fusion des canaux reducer /LastValueetc. Seuls les canauxis_available()(qui ont encore une valeur et n'ont pas été finish) entrent dansupdated_channels— aprèsfinish(), le numéro de version du canal avance mais ne déclenche plus aucun nœud.
prepare_next_tasks est appelé au début de tick pour décider quelles tâches exécuter à ce pas :
# 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, ..., for_execution=for_execution, ...
):
tasks.append(task)
# This section is an optimization that allows which nodes will be active
# during the next step.
if updated_channels and trigger_to_nodes:
triggered_nodes: set[str] = set()
for channel in updated_channels:
if node_ids := trigger_to_nodes.get(channel):
triggered_nodes.update(node_ids)
candidate_nodes: Iterable[str] = sorted(triggered_nodes)
elif not checkpoint["channel_versions"]:
candidate_nodes = ()
else:
candidate_nodes = processes.keys()
for name in candidate_nodes:
if task := prepare_single_task((PULL, name), None, ..., for_execution=for_execution, ...):
tasks.append(task)
return {t.id: t for t in tasks}Extrait de prepare_next_tasks 核心:441-513. Deux types de tâches s'y rejoignent :
- Tâche PUSH : on extrait les objets
Senddu canalTASKS; ce sont les fan-out produits parreturn Send(...)d'un nœud au pas précédent ; chaqueSendcorrespond à une tâche PUSH. - Tâche PULL : on utilise
updated_channelspour remonter viatrigger_to_nodes— ce mapping est calculé à la compilation « canal → liste des nœuds abonnés » ; l'union donne directement les nœuds candidats à déclencher à ce pas. Lesortedassure un ordre déterminé des tâches, sans effet sur le résultat d'exécution mais important pour le replay du checkpoint.
Chaque nœud candidat passe ensuite une seconde vérification _triggers (_triggers call:606-612) — « le canal a une nouvelle version et le nœud ne l'a pas encore vue » est la condition pour qu'une tâche soit réellement générée. Cette étape est la clé de la déduplication des nœuds PULL : même si updated_channels contient le canal trigger d'un nœud, si versions_seen[name] a déjà enregistré cette version, la tâche ne sera pas replanifiée.
La position de la couche algo dans la boucle Pregel :
Limites et échecs
- Si
versions_seen[INTERRUPT]n'est pas mis à jour, boucle infinie :should_interruptne renvoie une liste non vide que siany_updates_since_prev_interruptest vrai, or la mise à jour deversions_seen[INTERRUPT]se fait dans la logique implicite de_put_checkpointaprèsapply_writes. Si, lors d'un resume après interruption, la version seen du canal INTERRUPT n'est pas avancée correctement, la même mise à jour sera répétément considérée comme « nouvelle », conduisant à une boucle infinie interrupt-resume-interrupt (versions_seen INTERRUPT:163-168). bump_stepest le seul signal d'avance du step : siany(t.triggers for t in tasks)est faux, pas de bump — c'est la sémantique « null task ne fait qu'écrire dans les canaux sans avancer le step » (bump_step:259). Le null task sert à injecter des input writes avant le tick ; il ne doit pas faire sauter le compteur de step.finish()n'est appelé que si « aucun nœud n'a été déclenché » :bump_step and updated_channels.isdisjoint(trigger_to_nodes)(finish condition:336) — tant qu'un canal deupdated_channelsest encore abonné à un nœud, on ne finish pas. Cela signifie que le finish n'est envoyé qu'au « pas où le graphe meurt », pas à chaque pas.- Écriture dans un canal inconnu : warning, pas d'erreur : si
pending_writes_by_channelcontient un canal absent dechannels, on ne fait quelogger.warning(unknown channel warn:310-313). C'est par tolérance — par exemple un nœud qui écrirait dans un canal déjà finish ne doit pas faire planter tout le graphe. - Le chemin
fresh=Truedelocal_readcopie le canal : en lecturefresh, on faitchannels[k].copy()pour chaque canal (fresh read copy:213-219) et on applique les writes de cette tâche sur la copie — ainsi un nœud d'arête conditionnelle (conditional edge) peut lire sa vue locale « ce que je viens d'écrire mais qui n'est pas encore passé parapply_writes», sans impacter l'état global vu par les autres tâches. - Algorithme de hash de l'id de tâche dans
prepare_single_task: sicheckpoint["v"] > 1, on utilise_xxhash_str, sinon_uuid5_str(task_id_func:550) ; c'est une compatibilité entre versions de checkpoint — les task ids générés par uuid5 dans les anciens checkpoints doivent pouvoir être rejoués.
Résumé
Les trois fonctions pures de _algo.py — should_interrupt / apply_writes / prepare_next_tasks — capturent toute la couche sémantique de Pregel : quand interrompre, comment fusionner les writes, qui exécuter au pas suivant. Elles n'ont ni état ni IO, et sont pilotées par PregelLoop.tick + after_tick pour faire tourner la boucle BSP. Pour le pilote de boucle, voir /pregel/loop ; pour la couche d'exécution, voir /pregel/runner ; pour l'assemblage global, voir /pregel/pregel ; pour la sémantique des méthodes update / consume / finish / is_available des canaux, voir /channel/base-channel.
Voir la documentation officielle : LangGraph docs · README.