BaseChannel : le canal comme primitive d'état
Responsabilités
Dans le runtime de LangGraph, le « canal (channel) » est l'unité minimale d'état (state). L'instance Pregel produite par la compilation d'un StateGraph ne détient pas directement le dict utilisateur : elle détient un ensemble d'instances de BaseChannel — un canal par clé — et chaque canal stocke sa valeur, fusionne les écritures, expose une lecture, et sait aussi sérialiser son état courant sous forme de point de contrôle (checkpoint). BaseChannel est la classe de base abstraite partagée par tous les types de canaux (BaseChannel:19).
Elle ne stocke aucune donnée elle-même (__slots__ ne contient que key et typ, voir __slots__:22-26). Elle fixe seulement quatre membres abstraits que les sous-classes doivent implémenter : les deux propriétés de type ValueType / UpdateType, et les trois méthodes update / get / from_checkpoint (abstract methods:60-99). S'y ajoutent checkpoint / is_available / consume / finish / copy, cinq méthodes avec une implémentation par défaut qui soit se rabat sur get(), soit renvoie False en no-op.
Autrement dit, BaseChannel fixe en une fois l'interface de « comment l'état est stocké, modifié et tracé ». Le moteur Pregel n'a qu'à appeler update / get / checkpoint, sans avoir à savoir que la fusion est faite par un LastValue, un Topic ou un BinaryOperatorAggregate. C'est ce qui découple totalement le moteur de la forme de l'état.
Motivation de conception
Pourquoi modéliser l'état comme des « canaux » plutôt que de passer un dict entre les nœuds ?
- Sémantique unifiée de fusion d'écritures : au sein d'un même superpas (superstep), plusieurs nœuds peuvent écrire concurrently sur la même clé. Un dict ne sait pas exprimer les différentes règles de fusion — « liste à étendre, scalaire à écraser, compteur à additionner » ; le canal réunit ces stratégies derrière un
update(values: Sequence[Update]) -> bool, et la sous-classe décide du fold. - Vue immutable intra-superpas : un nœud lit le snapshot pris en début de superpas ; les valeurs écrites par d'autres nœuds dans le même superpas ne deviennent visibles qu'au superpas suivant.
getest une lecture de la valeur courante,updaten'est appelé qu'aux frontières de superpas parapply_writes(apply_writes:317-323), ce qui élimine toute race. - Sérialisable : chaque canal sait produire une valeur sérialisable via
checkpoint()(checkpoint:49-58) et se reconstruire viafrom_checkpoint(from_checkpoint:60-65). C'est sur ces deux interfaces que la familleBaseCheckpointSavers'appuie pour écrire l'état du thread en SQLite/Postgres. - Hooks de cycle de vie :
consumeetfinish(consume / finish:101-121) donnent aux canaux spéciaux (par ex.Topic(accumulate=False)ouLastValueAfterFinish) une chance de modifier leur état à la frontière de superpas ou en fin d'exécution ; le no-op par défaut indique que la plupart des canaux n'ont pas besoin de cette sémantique. - Auto-description des types : les deux propriétés
ValueType/UpdateType(ValueType / UpdateType:28-36) permettent à la compilation comme à l'exécution de retrouver « ce que ce canal stocke et ce qu'il accepte en écriture ».StateGraphles utilise dans_is_field_channel:1862-1887pour reconnaître les déclarationsAnnotated[..., SomeChannel].
Fichiers clés
BaseChannel class:19-26— définition de la classe de base abstraite, paramètres génériquesValue / Update / Checkpoint,__slots__ne déclarant quekeyettyp.ValueType / UpdateType:28-36— deux propriétés abstraites où la sous-classe déclare le type stocké et le type d'écriture.checkpoint:49-58— implémentation par défaut : renvoieself.get(), ouMISSINGpour un canal vide.from_checkpoint:60-65— méthode abstraite : construire un canal équivalent depuis une valeur de checkpoint, à implémenter par la sous-classe.get:69-73— lecture abstraite, lèveEmptyChannelErrorsur un canal vide.is_available:75-85— implémentation par défaut viaget()+ capture deEmptyChannelError; généralement surchargée par un test plus rapide.update:89-99— écriture abstraite, appelée par Pregel une fois à la fin de chaque superpas, dans n'importe quel ordre ; renvoieTruesi le canal a été modifié.consume / finish:101-121— hooks de cycle de vie, no-op par défaut renvoyantFalse.apply_writes:315-345— point d'entrée unique où Pregel réécrit les writes du superpas dans les canaux ;update/consume/finishy sont tous ordonnancés._is_field_channel:1862-1887— à la compilation,StateGraphtraduitAnnotated[..., channel]en instance de canal viaisinstance(item, BaseChannel).
Flux de données
Les écritures sur les canaux ont lieu dans apply_writes : Pregel regroupe les writes de toutes les tâches du superpas par canal, puis appelle update (apply_writes:317-323) :
# 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
# unavailable channels can't trigger tasks, so don't add them
if channels[chan].is_available():
updated_channels.add(chan)
# Channels that weren't updated in this step are notified of a new step
if bump_step:
for chan in channels:
if channels[chan].is_available() and chan not in updated_channels:
if channels[chan].update(EMPTY_SEQ) and next_version is not None:
checkpoint["channel_versions"][chan] = next_version
if channels[chan].is_available():
updated_channels.add(chan)Notez le second bloc : les canaux non écrits par aucune tâche dans le superpas reçoivent tout de même un appel update(EMPTY_SEQ). C'est un contrat implicite de l'interface BaseChannel — la méthode update d'une sous-classe doit accepter une séquence vide, généralement en renvoyant False pour « pas de modification ». La sémantique « effacer à chaque superpas » de Topic(accumulate=False) s'appuie justement sur cet appel à vide : l'update vide lui fait perdre la valeur du superpas précédent, qui devient invisible au superpas suivant.
Position du cycle de vie d'un canal dans un tick de Pregel :
Limites et échecs
EmptyChannelError: levé parget()quand le canal n'a jamais été écrit (get:70-73). L'implémentation par défaut deis_availables'appuie sur la capture de cette exception ; une sous-classe qui peut décider à moindre coût devrait la surcharger.- Sentinelle
MISSING:checkpoint()renvoielanggraph._internal._typing.MISSING— pasNone— quandget()lèveEmptyChannelError(checkpoint fallback:55-58), de sorte queNonepuisse être une valeur légale sans être confondu avec « vide ». updaterenvoieFalse: un retourFalsesignifie que le canal n'a pas bougé, Pregel ne met donc pas à jourchannel_versionsetprepare_next_tasksne déclenchera aucun nœud depuis ce canal (update call:319). C'est la base du critère de convergence de Pregel.- Appel à séquence vide :
apply_writesappelle aussiupdate(EMPTY_SEQ)sur les canaux inchangés (EMPTY_SEQ update:329). Le contrat deBaseChannel.updateest qu'une entrée vide renvoieFalse, maisTopic(accumulate=False)en profite pour s'effacer et peut renvoyerTrue: une sous-classe ne peut donc pas supposer « entrée vide = ne rien faire ». consumeetfinishen no-op par défaut (consume / finish:101-121) : le retourFalsesignifie que la plupart des canaux ne participent pas à la gestion du cycle de vie.apply_writesappellefinishà la fin debump_step(finish call:338) ; seuls les canaux spéciaux commeLastValueAfterFinishl'utilisent pour la sémantique « visible seulement après la fin de l'exécution ».copyse rabat sur checkpoint par défaut :BaseChannel.copyfait par défautself.from_checkpoint(self.checkpoint())(copy:40-47). Si le checkpoint est coûteux pour une sous-classe (par ex. une deep-copy d'une grosse liste), elle devrait surchargercopyavec une implémentation moins chère — c'est ce que faitTopic.copy.
Résumé
BaseChannel fixe en une fois « comment l'état est stocké, fusionné, sérialisé », et le moteur Pregel n'a qu'à suivre la séquence update → is_available → get → checkpoint → finish pour fonctionner sur n'importe quelle implémentation de canal. Pour les implémentations concrètes, voir /channels/last-value et /channels/topic-binop ; le déroulé des lectures/écritures par Pregel est dans /pregel/algo ; la persistance de la valeur renvoyée par checkpoint() est détaillée dans /checkpoint/base-saver.
Voir la documentation officielle : documentation LangGraph · README