Résumé
Sur un projet, j’avais un processeur entre deux topics Kafka: consommer sur le premier, valider, publier sur le second. La transformation était la partie facile. La difficulté était d’attendre que le producteur soit prêt, car il pouvait être mis en pause ou déconnecté de temps en temps, sans perdre un seul message consommé pendant ce temps.
Les opérateurs de fusion de RxJS ont l’air interchangeables et ne le sont pas. Un takeUntil suivi d’un repeatWhen se réabonne quand le producteur revient, mais tout ce qui a été reçu pendant la pause est perdu. Un switchMap qui renvoie une source vide tant que le producteur est absent perd les mêmes données. Un combineLatest fait à peine mieux, avec en prime un piège: filtrer directement sur le producteur, et le consommateur ne tient plus compte des déconnexions suivantes.
La réponse tient en deux opérateurs jumeaux. windowToggle dit quand émettre et quand ignorer; bufferToggle dit quand mettre en tampon. Le flux du consommateur bascule entre les deux états selon la disponibilité du producteur, puis les deux résultats sont fusionnés et aplatis pour publier chaque valeur une par une. Les données sont consommées quand le producteur reprend, et gardées quand il s’arrête.
L’exemple complet tourne avec NestJS, node-rdkafka et RxJS, et le dépôt qui accompagne l’article permet de rejouer le scénario. Le cas est étroit; la leçon, c’est que choisir un opérateur de fusion est une décision sur ce qu’on accepte de perdre.
Idées clés
- Le problème d’un processeur Kafka n’est pas la transformation, mais l’attente du producteur sans perdre ce que le consommateur reçoit pendant ce temps.
- takeUntil avec repeatWhen, switchMap et combineLatest reprennent tous quand le producteur revient, et perdent tous les données reçues pendant la pause.
- windowToggle décide quand émettre, bufferToggle quand mettre en tampon; fusionnés puis aplatis, ils ne laissent rien tomber.
- Filtrer directement sur le producteur est un piège: une fois disponible, le consommateur ignore les déconnexions suivantes.
Pourquoi j’ai écrit cet article
Dans ce tutoriel, je voulais expliquer quelques opérateurs de fusion et le problème que j’avais rencontré peu avant avec Kafka et RxJS: le projet sur lequel je travaillais avait une fonction de traitement entre un topic entrant et un topic sortant.