Introduction
Apache Kafka est une plateforme de streaming distribuée devenue la norme de facto pour la construction de pipelines de données en temps réel et d'applications de streaming. Son architecture est conçue pour offrir un débit élevé, une faible latence, une tolérance aux pannes et une évolutivité horizontale. Comprendre les composants essentiels de Kafka et le flux de données entre eux est crucial pour toute personne chargée d'exploiter, de surveiller ou de dépanner des clusters Kafka en production.
Ce guide explique l'architecture Kafka du point de vue d'un opérateur technique. Chaque section associe un concept fondamental à des commandes pratiques et des exemples de configuration. Vous apprendrez à inventorier votre environnement Kafka, à modifier les configurations en toute sécurité, à vérifier l'état du système, à diagnostiquer les pannes courantes et à suivre une liste de contrôle opérationnelle qui privilégie la sécurité et la réversibilité.
Nous utiliserons une entreprise fictive, StreamCo, pour illustrer des scénarios réels. StreamCo gère une plateforme de commerce électronique et utilise Kafka pour traiter les événements de commande, les mises à jour d'inventaire et les notifications clients. Leur cluster Kafka est en version 3.4.0, déployé dans trois centres de données avec trois brokers par centre.
Avant d'entrer dans les détails, établissons les composants fondamentaux de l'architecture Kafka.
Composants essentiels de l'architecture Kafka
L'architecture de Kafka repose sur quelques composants clés qui fonctionnent ensemble pour fournir un système de messagerie évolutif et tolérant aux pannes.
Brokers et clusters
Un broker Kafka est un processus serveur qui stocke les données et répond aux requêtes des clients. Un cluster Kafka se compose d'un ou plusieurs brokers. En environnement de production, on dispose généralement d'au moins trois brokers pour garantir la disponibilité et la tolérance aux pannes.
Chaque broker est identifié par un entier unique. Par exemple, dans l'environnement de StreamCo, les identifiants des brokers sont 0, 1 et 2 dans le centre de données principal.
Les brokers reçoivent les messages des producteurs, leur attribuent des offsets et les enregistrent sur disque. Ils répondent également aux requêtes des consommateurs en fournissant les messages à partir des segments de journal validés.
Commande pratique : Pour lister les brokers actifs d'un cluster, utilisez l'utilitaire kafka-broker-api-versions.sh fourni avec Kafka :
kafka-broker-api-versions.sh --bootstrap-server broker1.streamco.example:9092
Cette commande interroge les métadonnées du cluster et renvoie la liste des brokers avec les versions d'API prises en charge. C'est une opération en lecture seule qui ne modifie pas l'état du cluster.
Topics et partitions
Un topic est un canal logique dans lequel les producteurs publient des messages et à partir duquel les consommateurs lisent. Les topics sont divisés en partitions pour assurer l'évolutivité et le parallélisme. Chaque partition est une séquence ordonnée et immuable d'enregistrements, continuellement ajoutée à un journal de validation.
Les partitions permettent à Kafka de répartir les données sur plusieurs brokers. Par exemple, un topic commandes peut avoir 6 partitions. Ces partitions peuvent être réparties sur les trois brokers du centre de données principal, chaque broker hébergeant deux partitions.
Le nombre de partitions d'un topic est une décision de conception cruciale. Un plus grand nombre de partitions augmente le parallélisme en production et en consommation, mais accroît la charge sur le cluster. Une règle courante consiste à définir le nombre de partitions en fonction du débit attendu et du nombre de consommateurs dans un groupe de consommateurs.
Commande pratique : Pour créer un topic avec un nombre spécifié de partitions et un facteur de réplication, utilisez kafka-topics.sh :
kafka-topics.sh --create \
--bootstrap-server broker1.streamco.example:9092 \
--topic commandes \
--partitions 6 \
--replication-factor 3
Cette commande crée un topic nommé commandes avec 6 partitions et un facteur de réplication de 3, ce qui signifie que chaque partition est répliquée sur trois brokers.
Producteurs et consommateurs
Un producteur est une application qui publie (écrit) des messages dans des topics Kafka. Un consommateur est une application qui s'abonne à des topics et traite les messages publiés.
Les producteurs peuvent choisir la partition à laquelle envoyer un message. Par défaut, si aucune clé n'est spécifiée, le producteur utilise une stratégie de répartition circulaire pour distribuer les messages uniformément sur toutes les partitions du topic. Si une clé est spécifiée, Kafka utilise le hachage de la clé pour déterminer la partition, garantissant que tous les messages ayant la même clé vont dans la même partition, ce qui préserve l'ordre par clé.
Les consommateurs peuvent former des groupes de consommateurs. Dans un groupe de consommateurs, chaque partition est consommée par exactement un consommateur. Cela permet une consommation parallèle et un équilibrage de charge. Si un consommateur tombe en panne, un autre consommateur du groupe reprend ses partitions.
Extrait de configuration pratique : Voici une configuration minimale d'un producteur Java utilisant la bibliothèque cliente Kafka :
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
Cet exemple utilise des sérialiseurs de chaînes pour les clés et les valeurs. En production, on peut utiliser Avro ou Protobuf pour la gestion des schémas.
Le rôle de ZooKeeper
Historiquement, Kafka s'appuyait sur Apache ZooKeeper pour la coordination du cluster, l'élection du leader et la gestion des métadonnées. ZooKeeper maintient la liste des brokers, les configurations des topics et les informations sur le leadership des partitions.
Depuis Kafka 3.3, un nouveau mode KRaft (Kafka Raft metadata mode) est disponible, remplaçant ZooKeeper par un protocole de consensus intégré. KRaft simplifie le déploiement et améliore l'évolutivité. Cependant, de nombreux clusters de production utilisent encore ZooKeeper.
Pour StreamCo, leur cluster Kafka 3.4 utilise toujours ZooKeeper, il est donc important de surveiller la santé de ZooKeeper.
Commande pratique : Pour vérifier l'état de ZooKeeper, utilisez la commande echo srvr via nc :
echo srvr | nc zookeeper1.streamco.example 2181
Cette commande renvoie une ligne comme Mode: follower ou Mode: leader, indiquant le rôle du nœud ZooKeeper. Un cluster avec un nombre impair de nœuds assure un quorum.
Flux de données dans Kafka : du producteur au consommateur
Comprendre le chemin emprunté par un message du producteur au consommateur est essentiel pour diagnostiquer les problèmes de latence et de débit.
- Le producteur envoie un message : Le producteur sérialise le message et l'envoie à la partition leader du topic approprié. Le leader est le broker qui gère toutes les lectures et écritures pour cette partition.
- Le leader ajoute au journal : Le leader ajoute le message à son journal local et lui attribue un offset séquentiel.
- Réplication vers les suiveurs : Le leader réplique le message vers les réplicas suiveurs (si le facteur de réplication est supérieur à 1). Une fois le message répliqué sur le nombre requis de réplicas (selon le paramètre
acks), le leader accuse réception au producteur. - Le consommateur lit : Les consommateurs interrogent le leader pour obtenir des messages. Ils suivent leur offset dans la partition et peuvent valider les offsets dans Kafka ou en externe.
Exemple avec les paramètres acks :
acks=0: Le producteur n'attend aucun accusé de réception. C'est le plus rapide, mais sans garantie de livraison.acks=1: Le producteur attend uniquement l'accusé de réception du leader. Le leader écrit dans son journal local mais n'attend pas les suiveurs. C'est la valeur par défaut.acks=all(ouacks=-1) : Le producteur attend que tous les réplicas synchronisés accusent réception. Cela offre la garantie de durabilité la plus forte.
Pour le traitement des commandes de StreamCo, ils définissent acks=all pour éviter toute perte de commande même en cas de panne d'un broker.
Inventaire de la version et de l'environnement
Avant de modifier un cluster Kafka, vous devez comprendre les versions exactes, la topologie de déploiement et l'état actuel. C'est la première étape pour des opérations sûres.
Identifier la version installée
Utilisez la commande suivante pour connaître la version du broker Kafka :
kafka-broker-api-versions.sh --bootstrap-server broker1.streamco.example:9092 | grep -i version
Pour StreamCo, cela renvoie quelque chose comme broker1:9092 (id: 0 rack: null) -> 3.4.0. Notez l'ID du broker et le rack (s'il est configuré). L'information de rack est utile pour garantir que les réplicas sont répartis sur différents domaines de défaillance.
De plus, la distribution Kafka inclut le script kafka-configs.sh qui permet de récupérer la configuration du broker. Par exemple, pour vérifier le paramètre log.retention.hours :
kafka-configs.sh --bootstrap-server broker1.streamco.example:9092 \
--entity-type brokers \
--entity-name 0 \
--describe
Cette commande renvoie toutes les configurations du broker 0. C'est une opération en lecture seule.
Comprendre la topologie de déploiement
Cartographiez votre cluster : combien de centres de données, combien de brokers par centre, combien de partitions par topic et le facteur de réplication. Pour StreamCo :
- Centre de données principal : 3 brokers (IDs 0, 1, 2)
- Centre de données secondaire : 3 brokers (IDs 3, 4, 5)
- Centre de données de reprise après sinistre : 3 brokers (IDs 6, 7, 8)
Ils utilisent une configuration de réplication multi-régions avec une affectation personnalisée des partitions pour garantir que les réplicas se trouvent dans des centres de données différents.
Capturer l'état actuel
Avant tout changement, capturez un instantané des métriques et configurations pertinentes. Par exemple, pour obtenir la distribution actuelle du leadership des partitions :
kafka-topics.sh --describe --bootstrap-server broker1.streamco.example:9092 --topic commandes
Extrait de sortie attendu :
Topic: commandes Partition: 0 Leader: 0 Replicas: 0,3,6 Isr: 0,3,6
Topic: commandes Partition: 1 Leader: 1 Replicas: 1,4,7 Isr: 1,4,7
...
Cela indique quel broker est le leader pour chaque partition, les réplicas et les réplicas synchronisés (ISR). Si un réplica n'est pas dans l'ISR, il est en retard ou hors ligne.
Chemin de configuration sûr
Modifier les configurations Kafka peut avoir un impact significatif. Suivez un chemin sûr : observez, apportez un changement minimal, vérifiez et prévoyez un plan de retour en arrière.
Observation d'abord
Commencez toujours par des commandes en lecture seule. Par exemple, pour vérifier le paramètre actuel de rétention des journaux pour un topic :
kafka-configs.sh --bootstrap-server broker1.streamco.example:9092 \
--entity-type topics \
--entity-name commandes \
--describe
Cette commande renvoie toutes les configurations au niveau du topic. Supposons que la sortie affiche retention.ms=604800000 (7 jours). StreamCo souhaite augmenter la rétention à 14 jours pour permettre une analyse hors ligne.
Changement minimal justifié
Utilisez une modification de configuration dynamique pour mettre à jour la rétention du topic sans redémarrer les brokers :
kafka-configs.sh --bootstrap-server broker1.streamco.example:9092 \
--entity-type topics \
--entity-name commandes \
--alter \
--add-config retention.ms=1209600000
Cette commande change la rétention à 14 jours et l'applique dynamiquement. Le changement ne nécessite pas de redémarrage du broker.
Vérification et retour en arrière
Après le changement, vérifiez que la nouvelle configuration est active :
kafka-configs.sh --bootstrap-server broker1.streamco.example:9092 \
--entity-type topics \
--entity-name commandes \
--describe | grep retention.ms
Sortie attendue : retention.ms=1209600000.
Si le changement entraîne une utilisation inattendue du disque, vous pouvez revenir à la valeur d'origine en répétant la commande --alter avec retention.ms=604800000.
Rayon d'impact : Ce changement n'affecte que le topic commandes. Les autres topics conservent leurs paramètres de rétention d'origine.
Vérification et diagnostics
Une vérification efficace implique de contrôler la santé du cluster, l'état des brokers et le retard des consommateurs.
Contrôle de santé du cluster
Utilisez la commande kafka-broker-api-versions.sh pour vérifier rapidement que tous les brokers sont joignables :
kafka-broker-api-versions.sh --bootstrap-server broker1.streamco.example:9092
Si un broker est en panne, vous verrez une erreur de connexion pour ce broker. De plus, vous pouvez utiliser les métriques JMX de Kafka pour surveiller la santé des brokers. Par exemple, la métrique kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions indique le nombre de partitions sous-répliquées. Idéalement, ce nombre doit être nul.
Utilisation de JMX avec jconsole :
jconsole broker1.streamco.example:9999
Naviguez vers les MBeans et développez kafka.server -> ReplicaManager -> UnderReplicatedPartitions.
Surveillance du retard des consommateurs
Le retard des consommateurs est la différence entre l'offset le plus récent d'une partition et l'offset actuel du consommateur. Un retard élevé indique que les consommateurs prennent du retard.
Utilisez l'outil kafka-consumer-groups.sh :
kafka-consumer-groups.sh --bootstrap-server broker1.streamco.example:9092 \
--group processeur-commandes \
--describe
Sortie :
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
processeur-commandes commandes 0 10234 11200 966 consumer-1-a1c /192.168.1.10 consumer-1
processeur-commandes commandes 1 10450 11500 1050 consumer-1-b2 /192.168.1.10 consumer-1
Des valeurs de retard supérieures à un seuil (par exemple 1000) peuvent déclencher des alertes. StreamCo surveille le retard avec un exportateur Prometheus personnalisé.
Modes de défaillance et reprise
Kafka est conçu pour la tolérance aux pannes, mais des pannes surviennent malgré tout. Voici les scénarios de panne courants et comment s'en remettre.
Panne d'un broker
Symptôme : Un broker se met hors ligne. Les partitions qui avaient un leader sur ce broker élisent de nouveaux leaders parmi l'ISR. Si l'ISR contient d'autres réplicas, il peut y avoir une brève indisponibilité pendant l'élection du leader.
Détection : Utilisez kafka-broker-api-versions.sh pour voir si un broker est injoignable. Consultez les journaux du broker pour rechercher des messages FATAL ou ERROR.
Reprise :
- Déterminez la cause (panne matérielle, partition réseau, etc.).
- Si le broker peut être redémarré, redémarrez-le et surveillez sa réintégration dans le cluster. Il rattrapera les données manquées.
- Si le broker ne peut pas être redémarré, remplacez-le. Ajoutez un nouveau broker avec un ID différent et réaffectez les partitions si nécessaire.
Pour StreamCo, lorsque le broker 2 du centre de données principal est tombé en panne en raison d'un problème de disque, ils ont utilisé la commande suivante pour voir l'état des partitions :
kafka-topics.sh --describe --bootstrap-server broker1.streamco.example:9092 --topic commandes
Ils ont remarqué que la partition 2 avait un leader -1 (aucun leader) car tous les réplicas se trouvaient sur le broker en panne. Le facteur de réplication étant de 3, les autres réplicas sur différents brokers auraient dû prendre le relais, mais dans ce cas, l'ISR était vide car les suiveurs se trouvaient sur le même rack et avaient perdu la connectivité. Après le redémarrage du broker, l'ISR s'est reconstitué et le leadership a été rétabli.
Partitions sous-répliquées
Symptôme : La métrique UnderReplicatedPartitions reste non nulle pendant une période prolongée. Cela signifie que certaines partitions ont moins de réplicas que le facteur de réplication configuré.
Cause : Un broker est en panne ou un suiveur est trop en retard.
Reprise :
- Si un broker est en panne, remettez-le en service.
- Si un suiveur est en retard, vérifiez la connectivité réseau et les E/S disque sur ce broker.
- Vous pouvez également utiliser l'outil
kafka-reassign-partitions.shpour déplacer des partitions vers des brokers sains.
Exemple de réaffectation :
- Créez un fichier JSON
reaffectation.jsonavec les nouvelles affectations :
{
"version": 1,
"partitions": [
{"topic": "commandes", "partition": 0, "replicas": [1,4,7]}
]
}
- Exécutez la réaffectation :
kafka-reassign-partitions.sh --bootstrap-server broker1.streamco.example:9092 \
--reassignment-json-file reaffectation.json \
--execute
- Vérifiez avec l'option
--verify.
Problèmes de rééquilibrage des groupes de consommateurs
Symptôme : Les consommateurs d'un groupe se rééquilibrent constamment, ce qui entraîne des retards de traitement.
Cause : Cela se produit souvent lorsque les consommateurs mettent trop de temps à traiter les messages, dépassant max.poll.interval.ms, ou lorsqu'il y a des changements fréquents d'appartenance.
Reprise :
- Augmentez
max.poll.interval.msetsession.timeout.msde manière appropriée. - Assurez-vous que le temps de traitement des consommateurs est inférieur à l'intervalle d'interrogation.
- Vérifiez les redémarrages fréquents des consommateurs en raison d'exceptions.
Exemple de correction de configuration du consommateur :
max.poll.interval.ms=600000
session.timeout.ms=45000
heartbeat.interval.ms=15000
Liste de contrôle des opérations
Utilisez cette liste de contrôle pour les opérations de routine et lors de modifications de votre cluster Kafka.
| Élément à vérifier | Responsable | Fréquence | Commande / Métrique | Résultat attendu |
|---|---|---|---|---|
| Vérifier que tous les brokers sont actifs | Priya Shah, administratrice Kafka | Quotidienne | kafka-broker-api-versions.sh --bootstrap-server broker1:9092 | Tous les IDs des brokers listés sans erreurs |
| Vérifier les partitions sous-répliquées | Priya Shah | Toutes les heures | JMX UnderReplicatedPartitions | 0 |
| Surveiller le retard des consommateurs | Alex Chen, ingénieur données | Toutes les heures | kafka-consumer-groups.sh --group processeur-commandes --describe | Retard < 1000 pour toutes les partitions |
| Examiner l'utilisation du disque | Priya Shah | Hebdomadaire | df -h sur les hôtes des brokers | Utilisation du disque < 70 % |
| Sauvegarder les données ZooKeeper | Ravi Patel, DevOps | Quotidienne | Utiliser l'export d'instantané ZooKeeper | Export réussi |
| Tester le basculement | Priya Shah et Ravi Patel | Trimestrielle | Arrêter manuellement un broker et observer | Le transfert de leadership se fait en 30 secondes ; aucune perte de données |
| Examiner les configurations des topics | Priya Shah | Mensuelle | kafka-configs.sh --entity-type topics --describe | Les configurations correspondent aux besoins métier |
Responsabilité du propriétaire : Priya Shah, en tant qu'administratrice Kafka, est la principale responsable de la santé du cluster. Elle examine les résultats de la liste de contrôle chaque semaine avec l'équipe d'ingénierie et ajuste les seuils de surveillance si nécessaire.
Fréquence de révision : La liste de contrôle elle-même est revue chaque trimestre pour intégrer les nouveaux apprentissages et les changements du système.
Pièges courants et comment les éviter
1. Nombre de partitions insuffisant
Problème : Créer des topics avec trop peu de partitions limite le parallélisme et le débit. Par exemple, si un topic n'a qu'une seule partition, un seul consommateur d'un groupe peut la consommer, même s'il y a de nombreux consommateurs.
Comment le reconnaître : Retard élevé des consommateurs même avec de nombreux consommateurs, ou plateau du débit des producteurs.
Comment l'éviter : Avant de créer un topic, estimez le débit et le parallélisme des consommateurs. Utilisez un nombre de partitions multiple du nombre de brokers et permettant une croissance future. Pour StreamCo, ils ont calculé que le topic commandes nécessitait 6 partitions pour traiter 10 000 messages/s avec 3 consommateurs, permettant un parallélisme 3x.
Reprise : Vous pouvez augmenter les partitions plus tard avec kafka-topics.sh --alter --partitions N, mais sachez que l'augmentation des partitions modifie le partitionnement par clé et peut affecter les garanties d'ordre pour les clés existantes.
2. Facteur de réplication mal configuré
Problème : Définir replication.factor=1 signifie que les données sont perdues en cas de panne du broker. C'est une erreur courante en développement qui se retrouve accidentellement en production.
Comment le reconnaître : Vérifiez avec kafka-topics.sh --describe et voyez ReplicationFactor: 1.
Comment l'éviter : Définissez toujours un facteur de réplication d'au moins 3 en production et assurez-vous que les brokers sont répartis sur différents racks ou zones de disponibilité.
Reprise : Si vous avez un topic avec RF=1, vous pouvez utiliser kafka-reassign-partitions.sh pour augmenter le facteur de réplication, mais vous devez vous assurer que le cluster dispose de suffisamment de brokers pour héberger les réplicas.
3. Ignorer la santé de ZooKeeper
Problème : Même avec KRaft, de nombreux clusters dépendent encore de ZooKeeper. Si ZooKeeper tombe en panne, les brokers Kafka peuvent perdre la coordination et ne pas élire de leaders.
Comment le reconnaître : Les brokers affichent des erreurs comme Unable to connect to zookeeper dans les journaux. Les brokers peuvent s'arrêter si la connexion est perdue trop longtemps.
Comment l'éviter : Surveillez ZooKeeper avec echo srvr | nc zookeeper1 2181 et configurez des alertes. Assurez-vous que le cluster ZooKeeper a un nombre impair de nœuds (3 ou 5) pour le quorum.
Reprise : Redémarrez les nœuds ZooKeeper un par un pour rétablir le quorum. Redémarrez ensuite les brokers si nécessaire.
4. Ne pas surveiller le retard des consommateurs
Problème : Le retard des consommateurs peut augmenter silencieusement, entraînant des retards de traitement qui affectent les systèmes en aval.
Comment le reconnaître : Les applications en aval signalent des données obsolètes. Utilisez kafka-consumer-groups.sh pour voir des valeurs de retard élevées.
Comment l'éviter : Mettez en place une surveillance du retard avec Prometheus et Grafana. Alertez lorsque le retard dépasse un seuil pendant plus de quelques minutes.
Reprise : Augmentez le nombre de consommateurs en ajoutant des instances au groupe (si les partitions le permettent) ou optimisez la logique de traitement des consommateurs.
5. Modifications de configuration dynamique sans vérification
Problème : Modifier des configurations comme retention.ms sans vérifier l'effet peut entraîner une perte de données involontaire ou des erreurs de disque plein.
Comment le reconnaître : Après un changement, alertes inattendues ou rapports d'utilisateurs sur des données manquantes.
Comment l'éviter : Suivez toujours le chemin de configuration sûr : observez, modifiez de manière minimale, vérifiez et prévoyez un plan de retour en arrière.
Reprise : Revenez immédiatement sur la modification de configuration et enquêtez sur la cause première.
Conclusion
L'architecture de Kafka est robuste, mais elle nécessite une exploitation soigneuse pour maintenir la fiabilité et les performances. En comprenant les composants essentiels, le flux de données et les meilleures pratiques opérationnelles, vous pouvez gérer en toute sécurité les clusters Kafka en production.
Commencez par inventorier votre environnement et maîtrisez les commandes de diagnostic en lecture seule. Lorsque des modifications sont nécessaires, suivez un chemin sûr : observez, apportez un changement minimal, vérifiez et prévoyez un plan de retour en arrière. Utilisez la liste de contrôle des opérations pour attribuer les responsabilités et assurer une surveillance régulière. Évitez les pièges courants en concevant des topics avec un nombre de partitions et des facteurs de réplication appropriés, et en surveillant diligemment le retard des consommateurs et la santé du cluster.
Pour StreamCo, la mise en œuvre de ces pratiques a réduit les temps d'arrêt imprévus de 40 % et amélioré la capacité de l'équipe à répondre rapidement aux incidents. Appliquez ces principes à votre propre environnement Kafka pour obtenir des résultats similaires.
Comme prochaine étape, choisissez une vérification à faible risque de ce guide, comme le contrôle des partitions sous-répliquées, et intégrez-la à votre routine quotidienne. Progressez à partir de là.