Il ne suffit pas de s’appuyer sur l’ordre de Kafka au sein de la partition lorsque des milliers de sessions partagent cette même partition. Dans les systèmes de conversation et les tâches séquentielles, un message correctif comme « définissez la destination sur Paris » peut arriver après le message « réservez un voyage à Londres », mais son exécution avant le premier message détruirait le contexte. La solution décrite par Joshua Oluikpe dans InfoQ transfère une partie de la responsabilité de l’ordre à la couche applicative écrite en Go.
Ordre au sein de la session et parallélisme entre les sessions
La conception repose sur deux niveaux de workers. Le niveau de distribution reçoit les enregistrements Kafka et les achemine selon l’identifiant de session au moyen du hachage cohérent, de sorte que les messages d’une même session arrivent au même worker, tandis que les différentes sessions peuvent être réparties entre plusieurs workers. Ensuite, chaque session active reçoit une goroutine dédiée qui traite un seul message à la fois et dans l’ordre d’arrivée.
Ainsi, une session lente ne devient pas un obstacle pour les autres sessions, comme cela se produit dans un pool de workers plat lorsqu’une nouvelle tentative pour une seule session monopolise entièrement le worker. Les goroutines sont créées à la demande et supprimées après une période d’inactivité, ce qui fait dépendre la consommation de ressources du nombre de sessions actives plutôt que du nombre total de sessions possibles.
Effectuer de nouvelles tentatives sans dépasser les messages
La nouvelle tentative s’effectue au sein de la goroutine de la session elle-même. En cas d’erreur temporaire, telle que l’expiration du délai d’un service dépendant, le worker attend en appliquant un recul exponentiel avec une part d’aléatoire, puis retraite le même message avant de passer aux messages suivants. Comme les messages suivants restent dans le canal derrière le message en échec, le quatrième message ne peut pas dépasser le troisième.
Les erreurs irrécupérables, telles que des données non valides, sont quant à elles directement envoyées vers la liste des messages morts (DLQ). Cette distinction entre les erreurs temporaires et les erreurs définitives élimine le besoin d’une machine à états distincte pour gérer les nouvelles tentatives, mais signifie également que les messages d’une session donnée peuvent s’accumuler pendant la période de recul.
Validation sûre et récupération après un incident
Lorsque plusieurs sessions traitent leurs messages en parallèle, il n’est pas possible de valider uniquement l’offset terminé le plus élevé. Si l’offset 104 est terminé avant le 102 et que 104 est ensuite validé, un incident du consommateur peut entraîner le dépassement définitif de 102. Le système suit donc, pour chaque partition, les offsets en cours de traitement et ceux qui sont terminés, et ne déplace le point de validation que jusqu’au plus grand offset terminé de manière contiguë.
Le redémarrage peut par conséquent commencer à un offset plus ancien et retraiter certains messages. La conception utilise un identifiant d’événement stable pour faciliter la déduplication, mais ne prétend pas garantir une exécution « exactement une fois » ; les effets externes restent liés à la capacité des services dépendants à gérer les duplications de manière idempotente. Dans de rares cas, un trou bloqué peut être dépassé après l’expiration d’un délai défini afin de rétablir la progression de la partition, en l’enregistrant comme un échec opérationnel et en orientant le message vers la DLQ, ce qui constitue un compromis explicite entre continuité et exhaustivité du traitement.
Que faut-il pour une véritable mise en production ?
Les garanties d’ordre ne sont pas complètes avec la seule logique de base. Pendant la redistribution des partitions, le consommateur cesse de recevoir de nouveaux enregistrements, attend la fin du travail en cours, puis valide le dernier watermark contigu avant de transférer la propriété. Lorsque les canaux des workers sont pleins, il suspend temporairement la récupération des messages au lieu de les supprimer, tout en stockant les enregistrements déjà récupérés dans des buffers propres à chaque partition afin que des messages plus récents ne les dépassent pas.
Le système doit également détecter individuellement les offsets bloqués, assurer le suivi entre les frontières asynchrones de Kafka, maintenir les signaux de retard nécessaires à la mise à l’échelle automatique pendant la suspension de la récupération, et vider les travaux ainsi qu’écrire la DLQ avant l’arrêt. Le produit applique aussi le hachage cohérent du côté de l’envoi afin de garantir l’ordre des messages lors d’une production asynchrone.
Résultats et limites
Lors d’un test synthétique portant sur 50 000 messages et 10 sessions, avec injection d’erreurs temporaires dans 10 % des messages, le système a enregistré 14 027 messages par seconde sans aucune violation de l’ordre, avec une latence médiane de 11,7 millisecondes. Lors d’un test de charge sur 1 000 sessions, le débit effectif a atteint 48 805 messages par seconde sur un objectif de 50 000, sans erreur d’envoi, mais la latence de queue a dépassé le budget de 200 millisecondes en raison de la limite de débit de l’API Gateway, et non du mécanisme d’ordonnancement des sessions.
L’auteur indique que le système a traité plus de 40 millions de messages en production sans violation notable de l’ordre, avec seulement 460 messages orientés vers la DLQ, soit environ 0,002 %. Toutefois, ce chiffre repose sur la surveillance opérationnelle et non sur une vérification exhaustive de la séquence de chaque message, et le test de charge supérieur n’a pas vérifié indépendamment l’ordre de chaque session. La valeur pratique de la conception réside donc dans la transformation de l’ordre, qui n’est plus une hypothèse au niveau de Kafka mais une propriété imposée par l’application, tout en conservant la nécessité de tests explicites de séquencement et de garanties d’idempotence dans les services dépendants.