Le KvBus — un pub/sub en mémoire ou sur Redis, sans changer une ligne

Dans une application Socle, les workers ne s’appellent pas entre eux. Celui qui produit un événement le publie sur un sujet ; celui qui en a besoin s’y abonne. Aucun des deux ne connaît l’autre. Le KvBus est le bus qui porte ces messages, et aussi un magasin de clés que les workers partagent. Cet article montre son API, un abonnement réel, et la seule règle qu’il ne faut jamais enfreindre.

Le problème

Quand un worker en appelle un autre directement, les deux sont liés. On ne peut plus démarrer l’un sans l’autre, ni en ajouter un troisième qui voudrait le même événement. Et le jour où le travail se répartit sur deux instances, l’appel direct ne traverse pas le réseau.

Deux fonctions, une interface

Des clés partagées

get, set avec une durée de vie en secondes, del, scan par préfixe. Et des opérations atomiques : incr, compareAndSet, setIfNotExists — de quoi poser un verrou ou dédoublonner un message.

Des sujets

pub(sujet, contenu) publie ; sub(sujet, traitement) s’abonne. Le producteur ne sait pas qui écoute, ni combien.

Le KvBus est un composant Spring injecté par le constructeur. Il délègue tout à une implémentation choisie au démarrage par la variable KV_IMPL : in_memory (par défaut) ou redis. Le code des workers ne change pas d’une ligne.

Un abonnement réel

Voici comment le worker d’extraction de mots-clés de KEW, notre service d’entités nommées, s’abonne à ses demandes.

@Override
@SuppressWarnings("unchecked")
protected void onInitialize() {
    eventQueue = new LinkedBlockingQueue<>(500);

    // Souscription au canal d'entrée EventBus (SPEC §7.1)
    kvBus.sub("keywords.extract.requested", payload ->
        eventQueue.offer((Map<String, Object>) payload)
    );

    logger.info("[worker:{}][step:initialize] Souscrit au canal keywords.extract.requested", getName());
}

@Override
protected Map<String, Object> pollEvent() throws InterruptedException {
    return eventQueue.take();
}

Le traitement abonné ne fait qu’une chose : poser le message dans une file bornée. C’est voulu. En mémoire, le bus appelle les abonnés dans le fil du producteur : un traitement long à cet endroit le bloquerait. Le vrai travail — détection de langue, entités, enrichissement — se fait plus loin, dans les quatre fils virtuels que le worker ouvre et qui tirent la file par take().

La règle absolue : sub() dans initialize(), nulle part ailleurs

initialize() est appelée par le MOP une fois, au démarrage, avant que le worker ne reçoive son premier cycle. doWork(), elle, revient à chaque cycle.

En mémoire, chaque appel à sub() ajoute un traitement à la liste des abonnés du sujet. Un sub() posé dans doWork() ajoute donc un abonné par cycle. Après cent cycles, un seul message publié déclenche cent traitements. Rien ne plante : on voit seulement le même travail fait cent fois, et la mémoire qui monte. Sur Redis, c’est pire : chaque abonnement en trop occupe une connexion du pool.

Le signe qui ne trompe pas : dans les statistiques du bus (kvBus.getStats()), le compteur d’opérations sub qui grimpe pendant que l’application tourne.

En mémoire ou sur Redis

En mémoire Redis
Portée une JVM toutes les instances qui partagent le même Redis
Livraison immédiate, dans le fil du producteur par le réseau, dans un fil du client Redis
Message publié sans abonné gardé, jusqu’à 1 024 par sujet, et remis au premier abonné perdu
Survit au redémarrage non selon la persistance de Redis
Configuration aucune REDIS_URL, REDIS_NAMESPACE

Sur Redis, le pub/sub est celui de Redis : au plus une livraison, sans rattrapage. Un message publié pendant qu’un worker redémarre est perdu. Le préfixe REDIS_NAMESPACE isole les clés de deux applications qui partagent le même Redis.

Quand il ne faut rien perdre

Pour ces cas, la version 5.8.6 du bus offre une seconde paire de méthodes : publish(sujet, contenu) et subscribe(sujet, groupe, traitement). Les messages sont conservés. Un groupe d’abonnés se partage le travail. Un traitement qui lève une exception voit son message revenir ; après cinq livraisons ratées, le message part dans une file des rejets que l’on consulte par peekDlq. Sur Redis, cette voie repose sur les streams et leurs groupes de consommateurs.

Ce que ça donne en vrai

Le MOP lui-même utilise le bus pour ses commandes d’exploitation : il s’abonne à admin.shutdown_request, supervisor.drain_request et worker.restart. Redémarrer un worker depuis l’extérieur, c’est publier un message sur ce dernier sujet.

À retenir

  • Publier sans savoir qui écoute : c’est ce qui garde les workers indépendants.
  • sub() dans initialize(), et nulle part ailleurs.
  • pub/sub au plus une fois ; publish/subscribe quand il ne faut rien perdre.

Un programme à construire sur TheSocle ?

Retour en haut

Mentions légales · Confidentialité · Contact

© 2026 LMVI Conseil — SARL, SIREN 949 417 620 · [email protected]