1 points par GN⁺ 2024-11-14 | 1 commentaires | Partager sur WhatsApp
  • 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 = false et read_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() ou consumer.subscribe(), puis lit les records avec consumer.poll()
    • un consumer group se répartit le traitement des records d’un ensemble de topics
  • 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 = all a été utilisé
    • avec Bufstream, acks = 0 peut reconnaître une écriture sans attendre le stockage, ce qui peut entraîner la perte d’écritures commités
    • acks = 1 et acks = all bloquent 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 = true a été utilisée
  • Configuration du consumer

    • comme plusieurs documentations indiquent que l’auto-commit peut entraîner des pertes de données, enable.auto.commit = false a été utilisé dans l’ensemble
    • en l’absence d’offset commité, le auto.offset.reset par défaut démarre au dernier offset, ce qui ne garantit pas une livraison at-least-once
    • auto.offset.reset = earliest a é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_uncommitted lise les valeurs d’une transaction abortée est classé comme lecture abortée (G1a)
    • la documentation Kafka indique que read_committed empê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

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 client
    • subscribe ou assign : modifie l’ensemble de topics ou de partitions sur lequel le consumer effectue des poll
    • txn, poll, send : exécute une séquence de micro-opérations poll ou send
  • Dans la charge de travail non transactionnelle, chaque send ou poll ne 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 poll renvoie 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 poll jusqu’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
  • La charge de travail Abort a été ajoutée pour suivre le comportement des offsets de poll aprè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 poll sur un record, elle l’aborte intentionnellement, puis l’offset de poll suivant 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 avec node ... being disconnected ou timed out waiting for a node assignment, et poll se 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 0 avait 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 0 avant 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
  • É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 141 de la clé 5 a été renvoyée comme ayant été écrite avec succès à l’offset 274, mais tous les consumer.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 ProducerFencedException peut aussi être levée en cas de timeout de transaction
    • Le client Java Kafka utilise une TimeoutException dédiée pour la plupart des timeouts, mais dans ce cas il lève une ProducerFencedException
    • 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
  • 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 jt1234 et avait envoyé committed = false à EndTxn pour 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 EndTxn envoyé 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’offset 0 au 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 = latest peut 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 Produce retardé 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

 
GN⁺ 2024-11-14
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

    • J’aimerais vraiment faire une analyse de Kafka :-)
  • 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 surprenant
    Je 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)

    • C’est un peu surprenant, et je suis d’accord pour dire que la documentation n’explique pas bien ce point
      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
    • Le point essentiel est que le fait qu’un message ait été remis avec succès au client Kafka ne signifie pas que l’application l’a traité
      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
    • Je comprends cette fonctionnalité comme destinée aux situations à hautes performances
      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
    • On peut dans une certaine mesure justifier l’existence de cette fonctionnalité. Elle a été conçue pour un consommateur synchrone monothread, avec grosso modo une boucle qui appelle poll puis traite durablement les messages
      Le 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 à poll
      Par 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

    • L’introduction dit qu’« en octobre 2024, Bufstream n’a été déployé qu’auprès de clients sélectionnés », donc un système d’honneur pourrait être possible
      Il y a bien sûr un risque d’abus, mais cela peut être un compromis valable pour attirer certains clients
    • Le programme est soit open source, soit il ne l’est pas
      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

    • Je ne l’ai pas encore examiné, mais vous n’êtes pas le premier à le demander
      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 ?

  • 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 ?

    • Kafka appelle effectivement cela des transactions
      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/

    • Exact. WarpStream ne prend pas non plus en charge les transactions
  • Erratum : « Transactions may observe none, part, or all » devrait, je pense, être « Consumers may observe none, part, or all »

    • Les deux sont corrects, mais j’ai écrit transactions par souci de clarté
      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 ?

    • Jepsen est un outil qui vous ferait pleurer si vous ne saviez pas qu’il teste la base de données que vous développez
      Des larmes de joie, bien sûr. Le simple fait d’attirer l’attention de Jepsen est déjà un accomplissement
    • C’est un clone de Kafka. Kafka est essentiellement une file durable