- Lors de la vérification de Bufstream 0.1.0 à 0.1.3, un système de streaming compatible Kafka, 2 problèmes de disponibilité propres à Bufstream et 3 problèmes de sûreté ont été découverts ; les 5 ont été corrigés dans la version 0.1.3
- Les tests s’appuyaient sur Java Kafka Client 3.8.0 et sur les tests Jepsen existants pour Kafka/Redpanda, avec des réglages privilégiant la sûreté comme
acks = all,enable.idempotence = true,enable.auto.commit = falseetread_committed - Les problèmes de Bufstream incluaient l’arrêt de consommateurs et de producteurs, une réponse incorrecte avec l’offset
0, la perte de commits de transactions, ainsi que des pertes d’écritures reconnues dues à un bug de filtrage de la taille des réponses de l’API fetch - L’enquête a également révélé, dans le client Java Kafka et le protocole de transactions Kafka, des problèmes tels que le blocage indéfini de
Consumer.close(), des offsets de consommateurs imprévisibles, ainsi que des lectures abortées, écritures perdues et transactions déchirées - Jepsen estime que le protocole de transactions Kafka ne garantit pas explicitement l’ordre des requêtes client ni les numéros de transaction, si bien que, lors de l’utilisation du client Java officiel, la sûreté transactionnelle de Kafka et des systèmes compatibles Kafka peut être compromise
Structure de Bufstream et périmètre de la vérification
- Kafka est un système de streaming fournissant des journaux append-only répliqués et partitionnés, tandis que Bufstream est une implémentation alternative à Kafka qui privilégie la gouvernance des données et l’efficacité des coûts dans les environnements cloud
- Bufstream fournit, comme Kafka, des topics et des partitions, et fonctionne avec les clients Kafka standard
- le producer ajoute des records avec
producer.send() - le consumer se lie à une partition avec
consumer.assign()ouconsumer.subscribe(), puis lit les records avecconsumer.poll() - un consumer group se répartit le traitement des records d’un ensemble de topics
- le producer ajoute des records avec
- En l’intégrant à Buf Schema Registry, il est possible d’inspecter les records Protocol Buffer afin de prendre en charge la validation des records, le contrôle d’accès au niveau des champs et la conversion de formats de données avec d’autres systèmes
- Contrairement à Kafka, qui utilise des disques locaux et son propre protocole de réplication, Bufstream écrit les données directement dans de l’object storage
- il vise une réduction des coûts en tirant parti de la structure tarifaire du trafic de réplication de l’object storage
- les nœuds Bufstream peuvent fonctionner comme des VM stateless à autoscaling
- Bufstream est composé de trois sous-systèmes
- agent : service stateless fournissant l’API Kafka
- object store : stocke les chunks de records et les fournit aux lecteurs
- coordination service : utilise actuellement etcd, et détermine quels chunks sont commités ainsi que l’ordre des records
- En octobre 2024, Bufstream n’était déployé qu’auprès de certains clients, et sa documentation le présentait comme un « drop-in replacement » d’Apache Kafka, compatible avec les transactions Kafka et les exactly-once semantics, mais formulait peu de garanties de sûreté concrètes
Configuration client et hypothèses transactionnelles
- Comme dans ses tests précédents de systèmes compatibles Kafka, Jepsen a ajusté la configuration client afin d’obtenir un comportement plus sûr
-
Configuration du producer
- le réglage par défaut
acks = alla été utilisé - avec Bufstream,
acks = 0peut reconnaître une écriture sans attendre le stockage, ce qui peut entraîner la perte d’écritures commités acks = 1etacks = allbloquent jusqu’à ce que Bufstream soit certain de la persistance durable- pour éviter les appends en double lors des nouvelles tentatives automatiques du producer Kafka, la valeur par défaut
enable.idempotence = truea été utilisée
- le réglage par défaut
-
Configuration du consumer
- comme plusieurs documentations indiquent que l’auto-commit peut entraîner des pertes de données,
enable.auto.commit = falsea été utilisé dans l’ensemble - en l’absence d’offset commité, le
auto.offset.resetpar défaut démarre au dernier offset, ce qui ne garantit pas une livraison at-least-once auto.offset.reset = earliesta été utilisé afin que le consumer puisse observer l’ensemble du journal- une transaction Kafka se compose de l’ensemble de records envoyés par le producer et d’une map des offsets maximaux par partition pollés par le consumer
- ce n’est que lorsque la transaction est commitée que les records envoyés sont durables et finissent par être visibles pour les consumers
read_committed, tandis que l’offset commité avance également au moins jusqu’à l’offset spécifié dans la transaction - si la transaction n’est pas commitée, l’offset commité n’avance pas, et la visibilité des écritures peut varier selon la configuration du consumer
- le fait qu’un consumer
read_uncommittedlise les valeurs d’une transaction abortée est classé comme lecture abortée (G1a) - la documentation Kafka indique que
read_committedempêche G1a et garantit dans une certaine mesure la propriété selon laquelle toutes les écritures d’une transaction sont visibles ou aucune ne l’est, mais les tests Jepsen de Kafka, Redpanda et Bufstream ont observé des cycles d’écriture (phénomène proche de G0) et certaines formes de G1c
- comme plusieurs documentations indiquent que l’auto-commit peut entraîner des pertes de données,
Conception des tests
- Jepsen a testé Bufstream de la version 0.1.0 à la version 0.1.3, ainsi que plusieurs builds release candidate
- Le harnais de test utilisait le Bufstream test harness, la bibliothèque de test Jepsen et Java Kafka Client 3.8.0
-
Environnement d’exécution
- 3 à 5 nœuds Debian Bookworm ont été utilisés à la fois dans des containers LXC et sur des VM EC2
- 1 nœud était consacré à etcd, 1 nœud à Minio, et les autres servaient d’agents Bufstream
- les producers, consumers et clients d’administration étaient initialisés avec un seul nœud dans
bootstrap_servers, mais la découverte intelligente côté client n’était pas empêchée
-
Principaux réglages de sûreté
- auto-commit désactivé
acks = all- 1 000 tentatives
- idempotence activée
- niveau d’isolation
read_committed auto_offset_reset = earliest- création automatique de topics côté serveur désactivée
- les injections de pannes comprenaient la mise en pause de processus (
SIGSTOP), le crash (SIGKILL), le décalage d’horloge (clock_settime) et la partition réseau (iptables) - Bufstream étant divisé en agent, object store et coordination service, Jepsen a créé de nouveaux outils permettant d’injecter des pannes uniquement dans des sous-systèmes précis
- par exemple, les combinaisons variaient dans le temps, en ne faisant crasher que les nœuds Bufstream ou en ne mettant en pause que le coordinateur etcd
Charge de travail Queue et charge de travail Abort
- La charge de travail Queue analyse la sûreté en fonction du modèle de données de Kafka
- Chaque processus logique exécute un producer, un consumer et un client admin
- Une clé numérique identifie une topic-partition donnée
- Les clés sont choisies selon une fréquence exponentielle, si bien que certaines clés sont souvent accédées, tandis que d’autres le sont rarement
- Trois opérations de base sont utilisées
crash: termine le processus logique et le remplace par un nouveau clientsubscribeouassign: modifie l’ensemble de topics ou de partitions sur lequel le consumer effectue despolltxn,poll,send: exécute une séquence de micro-opérationspollousend
- Dans la charge de travail non transactionnelle, chaque
sendoupollne contient exactement qu’une seule micro-opération - Dans la charge de travail transactionnelle, plusieurs micro-opérations sont encapsulées dans une transaction Kafka
- L’analyse construit un mapping offset-vers-valeur pour chaque clé, puis recherche les erreurs
- Si plusieurs valeurs apparaissent au même offset, il s’agit d’un offset incohérent
- Si une même valeur apparaît à plusieurs offsets, il s’agit d’une erreur de duplication
- Si un record reconnu n’est jamais observé, il est considéré comme perdu ou non vu
- Si un
pollrenvoie une valeur envoyée par une opération abortée, il s’agit d’une lecture abortée - L’analyse vérifie aussi si une transaction observe ses propres écritures
- Après le test principal, les pannes sont résolues et l’on passe à l’étape des lectures finales
- Chaque processus lit toutes les topic-partitions depuis l’offset 0 et effectue des
polljusqu’au plus haut offset écrit connu - Si les lectures finales expirent et qu’un record reconnu n’est toujours pas observé, il est classé comme non vu
- Chaque processus lit toutes les topic-partitions depuis l’offset 0 et effectue des
- La charge de travail Abort a été ajoutée pour suivre le comportement des offsets de
pollaprès un abort de transaction- Le topic est limité à une seule partition, un seul processus, un seul producer et un seul consumer
- Après qu’une transaction a effectué un
pollsur un record, elle l’aborte intentionnellement, puis l’offset depollsuivant est classé comme avance, retour en arrière, retour plus loin en arrière ou autre
Cinq problèmes découverts dans Bufstream
-
Consommateurs bloqués (#1)
- De la version 0.1.0 à la 0.1.3-rc.8, la phase de lecture finale se bloquait fréquemment
consumer.poll()renvoyait immédiatement un résultat vide, mais des milliers de records acquittés restaient dans le journal- Cet état durait de plusieurs dizaines de secondes à plus d’une heure
- Dans un test, 691 records acquittés ont été envoyés durant les 120 premières secondes, et au début des lectures finales, 40 n’avaient été observés par aucun poller
- Ensuite, pendant plus d’une heure,
consumer.poll()n’a renvoyé aucun résultat, ce qui a fait expirer le test par timeout - La cause était qu’un nœud Bufstream redémarré pouvait renvoyer une valeur mise en cache obsolète du last stable offset et du high watermark
- Certaines bibliothèques clientes en concluaient qu’il n’y avait pas de records plus loin et se bloquaient ; Bufstream a appliqué un correctif dans la 0.1.3-rc.6 pour rafraîchir le cache au démarrage
-
Producteurs et consommateurs bloqués (#2)
- Même en 0.1.3-rc.6, des problèmes d’écritures non vues continuaient d’être observés après des pauses, crashs ou partitions touchant le coordinateur, le stockage ou un nœud Bufstream
- Dans certains cas, après une pause du coordinateur, tous les nœuds Bufstream étaient en cours d’exécution, mais le client entrait dans un état où il expirait en attendant
InitProducerId - Dans d’autres cas,
listOffsetséchouait avecnode ... being disconnectedoutimed out waiting for a node assignment, etpollse terminait mais ne renvoyait aucun résultat - Tuer puis redémarrer le nœud Bufstream résolvait le problème
- La cause était liée aux leases etcd
- L’agent Bufstream utilise les leases etcd pour suivre les agents actifs
- À cause d’une courte pause ou partition, etcd supprimait des clés liées au lease d’un agent, mais la mise à jour de suppression pouvait ne pas parvenir à l’agent
- L’agent se retrouvait alors sans savoir qu’il avait perdu son lease
- L’équipe Bufstream a ajouté une logique de polling supplémentaire, et les écritures non vues ont été globalement résolues en 0.1.3-rc.8
-
Faux offsets à zéro (#3)
- De la version 0.1.0 à la 0.1.3-rc.2, une valeur envoyée pouvait se voir attribuer l’offset
0, puis apparaître ensuite à un offset réel plus élevé - Cela se produisait même lorsque l’offset
0avait déjà été attribué bien plus tôt - Seul l’émetteur observait l’offset 0, tandis que les pollers observaient un offset plus élevé
- Dans un test de 2 minutes avec un seul nœud Bufstream et des pauses du processus etcd, 6 écritures ont reçu l’offset
0avant d’apparaître à un offset plus élevé - La cause était l’absence d’un champ requis dans la réponse d’erreur de Bufstream
- Bufstream envoyait une requête de commit de journal à etcd et etcd la traitait, mais à cause d’une pause ou d’une partition, Bufstream pouvait expirer en attendant la réponse
- Bufstream envoyait un code d’erreur au client, mais ne définissait pas l’offset du record envoyé à
-1, le signal d’erreur - Le client Java Kafka interprétait cela comme une réponse de succès à l’offset
0 - Franz-go, utilisé par la suite de tests de Bufstream, interprétait ce message comme une erreur, si bien que le problème n’apparaissait pas dans les tests
- Bufstream a corrigé le problème en 0.1.3-rc.6, et Jepsen ne l’a plus observé ensuite
- De la version 0.1.0 à la 0.1.3-rc.2, une valeur envoyée pouvait se voir attribuer l’offset
-
Écritures transactionnelles perdues (#4)
- En 0.1.2, une perte d’écritures se produisait fréquemment : certains records de transactions committées disparaissaient et n’étaient plus jamais observés
- Dans un test, sur 100 secondes et 6 761 transactions d’écriture, 240 records écrits par des transactions committées ont été perdus
- Dans l’exemple, la valeur
141de la clé5a été renvoyée comme ayant été écrite avec succès à l’offset274, mais tous lesconsumer.poll()sautaient cet offset - La cause était un bug dans le mécanisme de sûreté de concurrence ajouté en 0.1.2
- Ce mécanisme attribue un numéro unique à chaque transaction au sein d’une epoch de producteur afin de pallier le manque d’idempotence du protocole de transaction Kafka
- À cause d’un bug dans la logique de suivi des numéros de transaction, certains commits étaient ignorés à tort lorsque plusieurs transactions étaient committées sur plusieurs epochs
- Une transaction semblant committée pouvait en réalité être abortée, ou inversement
- Jepsen a découvert ce bug grâce à un timeout de transaction fixé très bas, à 1 seconde
- Bufstream a identifié le problème quelques heures après la publication de la 0.1.2, a empêché la mise à niveau des clients, et aucun client n’a migré vers la 0.1.2
- Le correctif a été inclus dans la 0.1.3-rc2
-
Écritures perdues dues au filtrage côté serveur (#5)
- En 0.1.3-rc.8, une courte fenêtre de perte d’écritures apparaissait fréquemment après de petites défaillances, comme une pause du processus Bufstream ou du coordinateur, ou une partition entre les deux
- La perte de données se produisait indépendamment de l’utilisation ou non de transactions
- Dans un test de 5 minutes, 22 records sur 16 770 ont été acquittés, mais aucun consommateur n’a pu les récupérer via poll
- Certains records étaient visibles par les pollers pendant un temps, puis disparaissaient ensuite des résultats de poll
- La cause était la logique de limitation de la taille des réponses de l’API fetch ajoutée en 0.1.3-rc.8 pour contourner un bug d’une interface web Kafka populaire
- Un bug dans la logique de filtrage cachait des records aux consommateurs en retard, ce qui se manifestait comme une perte d’écritures
- Bufstream a corrigé le problème en 0.1.3-rc.12
Problèmes du client Java Kafka et du protocole Kafka
-
KIP-588 : une ProducerFencedException trompeuse
- Pendant les tests, l’erreur
ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one.est apparue fréquemment - Cette erreur s’est produite même dans des tests où tous les producers recevaient un ID transactionnel unique, ce qui a fait perdre du temps pour en comprendre la cause
- KIP-588 indique qu’une
ProducerFencedExceptionpeut aussi être levée en cas de timeout de transaction - Le client Java Kafka utilise une
TimeoutExceptiondédiée pour la plupart des timeouts, mais dans ce cas il lève uneProducerFencedException - Même lorsqu’il n’existe en réalité aucun producer concurrent, le message d’erreur affirme qu’une deuxième instance de producer existe
- KIP-588 est ouvert depuis deux ans, et Jepsen recommande à l’équipe Kafka de modifier le message d’erreur
- Pendant les tests, l’erreur
-
KAFKA-17734 : Consumer.close() peut bloquer indéfiniment
- Dans les tests de Bufstream comme de Kafka, les tests se bloquaient toutes les quelques heures à cause d’un bug du client Java
Consumer.close()bloque par défaut sur les E/S réseau- Le paramètre de timeout de
close()devrait empêcher un blocage indéfini, mais il ne fonctionnait pas - L’approche consistant à appeler
consumer.wakeup()depuis un thread séparé pour interrompre un consumer bloqué dans les E/S n’a pas non plus été efficace - Jepsen estime qu’un programme de longue durée doit pouvoir libérer des ressources comme les clients, connexions, threads et mémoire dans un délai raisonnable, même en cas d’erreur réseau, et a ouvert KAFKA-17734
-
KAFKA-17582 : offset de consumer imprévisible après un échec de transaction
- La documentation officielle de Kafka dit très peu de choses sur ce que doit devenir l’offset du consumer lorsqu’un commit de transaction échoue
- La documentation de conception Kafka de Confluent indique que, si une transaction est abortée, la position du consumer revient à sa valeur précédente, mais le client Java réel ne se comporte pas toujours ainsi
- Les résultats du workload d’abort ont montré que, même sur un cluster sain, le comportement après un abort était difficile à prédire
- La plupart des paires de transactions avancent vers un offset plus lointain
- Certaines reviennent en arrière vers un offset précédent
- Tous les retours en arrière étaient liés à un événement de rebalance, et aucune avancée n’impliquait de rebalance
- D’après la réponse côté Kafka, ce comportement est intentionnel
- Le consumer continue d’avancer
- Si un rebalance se produit, il peut revenir à un point arbitraire selon l’offset committé
- Les utilisateurs doivent repositionner manuellement la position du consumer en cas d’abort de transaction
- Jepsen a ouvert KAFKA-17582 et proposé de documenter ce comportement, ainsi que d’envisager un retour arrière par défaut lors d’un abort de transaction
- Le workload Queue a lui aussi été modifié pour repositionner explicitement le consumer en arrière
-
KAFKA-17754 : perte d’écriture, lecture abortée, transaction déchirée
- Dans Bufstream 0.1.0 à 0.1.3, de simples pauses de processus Bufstream, pauses du coordinateur, crashs et partitions réseau ont permis d’observer des lectures abortées, des écritures perdues et des violations d’atomicité
- L’analyse a mis en évidence un défaut fondamental du protocole transactionnel de Kafka
- Dans l’exemple, le client exécutait une transaction avec l’ID transactionnel unique
jt1234et avait envoyécommitted = falseàEndTxnpour l’aborter, mais 15 appels àpoll()ont observé les écritures de la transaction abortée - D’autres écritures de la même transaction n’ont été observées par aucun poller
- L’examen conjoint de la capture de paquets et des logs Bufstream a montré que la cause était un message de commit retardé
- Un commit
EndTxnenvoyé plusieurs transactions plus tôt a été traité en retard sur un nœud - Le client avait déjà poursuivi avec les transactions suivantes
- Le commit retardé a été appliqué à la transaction courante, ce qui a committé seulement le début de la transaction, tandis que le reste a été traité comme une transaction distincte et aborté
- Le protocole Kafka est conçu pour permettre à un client d’envoyer des requêtes sur plusieurs connexions TCP et à plusieurs nœuds, mais il ne possède pas de numéro de séquence définissant l’ordre des requêtes d’un même client
- Il n’existe pas non plus de notion de numéro de transaction, si bien que le serveur ne peut pas savoir quelle transaction le client cherchait à terminer lorsqu’il reçoit un message de commit ou d’abort
- Les situations suivantes deviennent donc possibles
- Une transaction qui semblait committée est en réalité abortée
- Une transaction abortée est en réalité committée
- Une transaction déchirée se produit, où seule une partie des écritures de la transaction est conservée et le reste est perdu
- Le client Java Kafka officiel traite les timeouts comme retryables et peut envoyer automatiquement plusieurs messages
EndTxn, de sorte que le problème peut survenir même si l’utilisateur n’appelle commit ou abort qu’une seule fois par transaction - Jepsen a aussi observé dans Kafka des lectures abortées et des transactions déchirées via des pauses de processus, et a ouvert KAFKA-17754
- Les ingénieurs Kafka estiment que KIP-890 pourrait corriger ce problème
- KIP-890 modifie le protocole transactionnel en incrémentant l’epoch du producer à chaque transaction
- Comme le serveur rejette les messages d’un ancien epoch, cela peut empêcher qu’un message de commit d’une transaction passée se retrouve dans une transaction ultérieure
- Dans la version 0.1.3, Bufstream a ajouté un mécanisme utilisant la révision etcd comme horloge logique pour réduire la fréquence du problème, mais cela n’empêche pas les réordonnancements entre le client et Bufstream
- Jepsen a continué d’observer des lectures abortées, des écritures perdues et des transactions déchirées en 0.1.3, et estime qu’une correction côté client est nécessaire
Résumé global des résultats
- Les 5 problèmes propres à Bufstream ont tous été corrigés
- #1 : un « highest stable offset » en retard bloquait le consumer, pas besoin de panne ; corrigé dans 0.1.3-rc.6
- #2 : l’expiration de lease etcd bloquait le producer/consumer, pause nécessaire ; corrigé dans 0.1.3-rc.8
- #3 : offsets zéro parasites, pause nécessaire ; corrigé dans 0.1.3-rc.6
- #4 : écritures transactionnelles perdues, pas besoin de panne ; corrigé dans 0.1.3-rc.2
- #5 : écritures perdues à cause du filtrage côté serveur, pause nécessaire ; corrigé dans 0.1.3-rc.12
- Les problèmes liés à Kafka restent présents
- KIP-588 : message d’erreur incorrect lors d’un timeout de transaction, non résolu
- KAFKA-17734 :
ConsumerClient.close()peut se bloquer indéfiniment, non résolu - KAFKA-17582 : après un échec de transaction, l’offset consumer est imprévisible, non résolu
- KAFKA-17754 : perte d’écriture, lecture avortée, transaction déchirée, non résolu
- Jepsen rappelle que la vérification expérimentale de sûreté peut prouver l’existence de bugs, mais pas leur absence
- En particulier, à cause de KAFKA-17754, il estime difficile de déterminer s’il existe d’autres cas de perte d’écriture dans Bufstream
Recommandations pour les utilisateurs et l’exploitation de Bufstream
- Les utilisateurs qui utilisent les transactions Bufstream avec le client Java Kafka officiel doivent considérer que les transactions peuvent ne pas être sûres actuellement
- Une transaction avortée peut en réalité être commitée
- Une transaction commitée peut en réalité être avortée
- Une transaction peut être coupée en deux, avec seulement une partie de ses effets conservée
- Bufstream estime que le client Franz-go est moins vulnérable à ce problème, mais Jepsen n’a pas testé Franz-go avec les mêmes techniques que dans ce travail
- D’autres clients peuvent être vulnérables ou non
- Les utilisateurs de versions antérieures à Bufstream 0.1.3 peuvent rencontrer les problèmes suivants
producer.send()renvoie à tort l’offset0au lieu de l’offset réel- Problème de disponibilité métastable où le client reste bloqué
- Jepsen recommande la mise à niveau vers 0.1.3
- Il juge que l’architecture globale de Bufstream semble saine
- Déterminer l’ordre de chunks de données immuables via un service de coordination comme etcd est une approche relativement simple, déjà éprouvée dans les systèmes OLTP et de streaming
- Côté exploitation, deux améliorations sont recommandées
- Si, au démarrage, une requête de fichier partagé vers le stockage échoue, le cluster peut crasher ; Jepsen a recommandé d’ajouter des retries, et Bufstream a ajouté une couche de retry
- Lorsque des dépendances sont indisponibles, il est recommandé que l’agent continue à s’exécuter, fournisse du backpressure et un état du système, et récupère plus en douceur, plutôt que de mourir immédiatement
- À partir de la version 0.1.3, Bufstream a ajouté de la logique de retry supplémentaire pour etcd, mais un supervision constante reste nécessaire pour maintenir l’état online
- Les utilisateurs doivent vérifier qu’un process supervisor est en place et qu’il continue de fonctionner sans abandonner pendant les longues pannes
Nécessité de documenter les transactions Kafka et de corriger le protocole
- La documentation officielle de Kafka parle très peu des transactions, si bien que les utilisateurs doivent combiner plusieurs sources ambiguës et contradictoires
- Jepsen a recommandé à l’équipe Kafka de créer un document central clarifiant la sémantique des transactions, et mentionne KAFKA-17671
- Ce document devrait au minimum préciser les points suivants
- quand un consumer observe des offsets croissants de façon monotone
- quand un consumer peut sauter un record acquitté
- si un rebalance peut avoir un effet au milieu d’une transaction
- quand les offsets d’écriture du producer augmentent de façon monotone
- quand G0, G1a, G1b, G1c, les lectures fracturées et la lecture de ses propres écritures transactionnelles sont légales
- quelle signification ont les valeurs de retour de
poll()et les offsets après une transaction avortée - comment traiter les erreurs de transaction, les erreurs pendant un abort et les erreurs pendant un rewind
- La documentation Confluent répète que les valeurs par défaut de Kafka fournissent une livraison at-least-once, mais Jepsen souligne que cela ne semble pas être vrai
auto.offset.reset = latestpeut faire apparaître des records non traités comme « committed »- La documentation Confluent sur la gestion des offsets indique elle aussi un risque de perte de progression des messages en cas de crash avec l’auto-commit par défaut
- La documentation affirmant que le consumer est rembobiné lors d’un abort de transaction ne correspond pas non plus au comportement réel
- Jepsen estime que le protocole de transactions Kafka doit être corrigé en profondeur
- Le protocole suppose implicitement une livraison ordonnée et fiable, alors qu’il existe des pauses de processus, une fiabilité réseau imparfaite, une latence non nulle et une livraison non ordonnée entre plusieurs sockets TCP
- Le protocole Kafka répartit les messages entre plusieurs nœuds et sockets TCP, et le client retry automatiquement les messages
- Il n’existe pas de numéro de séquence pour restaurer l’ordre des messages d’un même client, ni de numéro de transaction pour vérifier la cible d’une transaction
- KIP-890 cherche à garantir un ordre plus strict en incrémentant l’epoch à chaque commit de transaction
- Les bibliothèques clientes peuvent aussi aider en réinitialisant le producer pour incrémenter l’epoch lorsqu’un message n’est pas acquitté
- Java Kafka Client 3.8.0 est vulnérable à ce problème
- Jepsen estime que Franz-go peut atténuer ou éviter le problème en effectuant une réinitialisation lors d’un timeout, mais n’a pas examiné les autres bibliothèques clientes
Travaux futurs
- De nombreux utilisateurs s’appuient sur les « exactly-once semantics » de l’API Kafka Streams plutôt que de manipuler directement les transactions ; il serait donc possible d’examiner à l’avenir la correction des applications Streams
- En enquêtant sur KAFKA-17754, Jepsen a également rencontré des écritures non vues dans Kafka, mais n’a pas pu les analyser faute de temps
- Une écriture non vue peut être le signe d’une transaction suspendue, d’un consommateur bloqué ou d’une perte de données
- La question reste ouverte de savoir si un message
Produceretardé pourrait entrer dans une transaction future et violer les garanties transactionnelles - Jepsen soupçonne aussi que le client Java Kafka puisse réutiliser un numéro de séquence lors d’un timeout de requête, avec la possibilité qu’une écriture ait été acknowledge mais soit silencieusement discardée
- Lorsqu’un événement de rebalance survient, la position d’un consumer peut avancer ou reculer, mais les règles exactes ne sont pas claires
- Si Kafka documente le comportement attendu, Jepsen souhaite le vérifier
- Jepsen explique qu’étant un processus aléatoire, il est difficile d’explorer les anomalies rares
- Un problème qui ne se produit qu’une seule fois est très difficile à déboguer et à reproduire
- Bufstream utilise également Antithesis, qui exécute l’ensemble d’un système distribué dans un hyperviseur déterministe et un réseau simulé
- Combiner la génération de workload et la vérification d’historique de Jepsen avec l’environnement déterministe et rejouable d’Antithesis pourrait améliorer la reproductibilité des tests
1 commentaires
Avis sur Hacker News
Si, en enquêtant sur des problèmes comme KAFKA-17754, on découvre aussi des écritures invisibles dans Kafka, il semble temps que Jepsen se replonge sérieusement dans Kafka
La dernière enquête remontait à 2013 (https://aphyr.com/posts/293-call-me-maybe-kafka, Kafka 0.8 bêta), et aujourd’hui on a l’impression d’être au stade où l’on commence tout juste à découvrir plusieurs problèmes dans Kafka lui-même
Des choses comme « une écriture peut être confirmée puis discrètement jetée » font assez peur
Le fait qu’avec la valeur par défaut
enable.auto.commit=true, un consommateur Kafka puisse committer des offsets indépendamment du fait que l’application les ait réellement traités est très surprenantJe n’avais jamais compris l’auto-commit ainsi, et une telle valeur par défaut me paraît absurde
La documentation n’est pas parfaitement claire, mais dans l’ensemble je l’avais lue comme disant que les offsets ne sont committés qu’une fois le traitement terminé
Je comprenais le réglage de l’intervalle d’auto-commit comme un moyen de réduire la fenêtre de retraitement en double, et non de perdre des messages, comme on s’y attendrait avec une sémantique au moins une fois (at-least-once)
Sans commit explicite, Kafka n’a aucun moyen de savoir si le message a été traité
Kafka suppose que le message qu’il a transmis a été traité immédiatement
L’auto-commit, c’est un peu comme tendre un cornet de glace puis se retourner aussitôt en supposant que la personne l’a mangé. Certaines personnes peuvent le faire tomber dès qu’elles le reçoivent et ne pas en manger une seule bouchée
Si l’on veut cette garantie, il faut envoyer un accusé de réception explicite
Par exemple, si tout ce que vous faites est d’écrire le message dans une base de données, le message est considéré comme acquitté dès qu’il arrive dans le callback du handler côté client
Mais en réalité, vous voulez probablement qu’il ne soit acquitté qu’après la réussite de l’insertion en base
Si la base de données devient inaccessible à cause du réseau, de Kubernetes, d’une configuration de pare-feu, etc., et qu’un ingénieur tente de redémarrer pendant ce temps, faisant tomber le client, il est facile de se retrouver avec des messages non traités
Un autre système peut déterminer s’il y a eu échec, et cette fonctionnalité permet de déplacer la borne supérieure pour réduire le retraitement
Cela dit, si le timing s’y prête et qu’une panne survient, il faut supposer qu’après redémarrage on peut recevoir à nouveau une partie déjà traitée
Le problème, c’est lorsqu’aucun traitement de ce type n’a lieu avant l’auto-commit
À la lecture, cela semble conçu pour committer longtemps après le traitement, mais le fait que ce soit automatique tout en ne devant committer que les éléments situés quelques millisecondes avant le moment de l’auto-commit paraît presque contradictoire
pollpuis traite durablement les messagesLe point qui prête à confusion est que la vérification de l’auto-commit ne se produit pas de façon asynchrone après un timeout, mais au moment de l’appel suivant à
pollPar conséquent, on ne devrait pouvoir perdre des écritures que dans les cas où l’on stocke simplement les messages sans les traiter durablement avant de rappeler
poll, par exemple avec du traitement asynchrone, des délais, des files, etc.Cela repose sur le comportement documenté de la bibliothèque cliente Java (https://kafka.apache.org/32/javadoc/org/apache/kafka/clients...) ; savoir si l’implémentation actuelle fait vraiment cela est une autre question
Le protocole Kafka est coincé entre haut niveau et bas niveau, et ne fait aucun des deux particulièrement bien
L’auto-commit est une fonctionnalité de haut niveau qui aide à construire facilement des applications simples, mais si on ne l’utilise pas de la manière attendue, elle peut évidemment échouer
Aujourd’hui, je pense que les utilisateurs finaux devraient utiliser des implémentations de plus haut niveau qui gèrent correctement les détails plutôt que d’utiliser directement les clients Kafka. Pour les usages data, cela veut dire un moteur de stream processing ; pour les usages applicatifs, quelque chose comme un moteur d’exécution durable
En regardant la page produit (https://buf.build/product/bufstream), je me demande comment concilier l’affirmation selon laquelle il « s’exécute uniquement dans votre VPC AWS ou GCP et ne contacte pas l’extérieur » avec une facturation à l’usage de « 0,002 $ par Gio avant compression »
J’imagine mal toute l’activité reposer sur un système d’honneur
Il y a bien sûr un risque d’abus, mais cela peut être un compromis valable pour attirer certains clients
Si le code source n’est pas publié, il ne faut jamais croire l’affirmation selon laquelle il « ne contacte pas l’extérieur »
« Le protocole de transactions de Kafka est fondamentalement cassé et doit être révisé », ça fait mal à entendre
Mais comme toujours, l’enquête et l’article sont excellents
Je me demande si Kyle a déjà examiné NATS JetStream. Je serais curieux de savoir ce qu’il en penserait
Plusieurs personnes ont suggéré que ce serait… comment dire… intéressant :-)
Je n’arrive pas à trouver le projet GitHub de bufstream ; où se trouve-t-il ?
Curieusement, celui-ci n’a pas de licence non plus
En lisant les articles de blog et la documentation associés, il semble que la « livraison exactement une fois » de Kafka soit définie comme une propriété d’une opération lire-traiter-écrire, où un worker lit dans le topic 1 et écrit dans le topic 2, les deux topics appartenant au même système Kafka logique
Si c’est bien cela, ne vaudrait-il pas mieux appeler ça une transaction ?
Mais il y a deux façons de voir le « exactement une fois »
L’une signifie que, comme dans une transaction de base de données, les effets ne doivent pas être dupliqués ni disparaître
L’autre ressemble davantage à une propriété de graphe de flux de données sur les relations entre messages à travers des topics-partitions, et se rapproche un peu plus de la cohérence dans ACID
De la même manière qu’un système transactionnel sérialisable garantit une cohérence au niveau du domaine, on peut utiliser des transactions pour atteindre cette propriété de flux de données
Par exemple, la sérialisabilité garantit que les invariants préservés par chaque transaction prise isolément le sont aussi dans l’historique concurrent
On peut voir cela comme la manière dont Kafka tente d’atteindre une « sémantique exactement une fois »
À ne pas confondre avec https://www.warpstream.com/
Erratum : « Transactions may observe none, part, or all » devrait, je pense, être « Consumers may observe none, part, or all »
La sémantique des consommateurs hors transaction est plus floue
Toutes les lectures de cette charge de travail se produisent dans un contexte transactionnel et passent par le chemin de commit des offsets transactionnels
Je me demande à quoi sert ce logiciel. Instrumentation ? Boîte noire ?
Des larmes de joie, bien sûr. Le simple fait d’attirer l’attention de Jepsen est déjà un accomplissement