La sauvegarde opérationnelle de Kafka exige d'aligner une stratégie de récupération sur des objectifs RPO et RTO explicites, la topologie du cluster (KRaft vs. ZooKeeper) et la posture de sécurité. Ce guide couvre cinq méthodes de production — MirrorMaker 2, Kafka Connect Replicator, instantanés de volume, export de stockage hiérarchisé et kafka-dump-log — avec des commandes selon la version, des signaux de vérification, des diagnostics de modes de défaillance et des procédures de retour arrière. Tous les exemples utilisent des espaces réservés $UPPER_SNAKE_CASE ; remplacez-les par vos valeurs avant exécution. Testez chaque procédure dans un environnement non-productif d'abord ; les services gérés (Confluent Cloud, Amazon MSK, Aiven, Redpanda) peuvent restreindre l'accès au niveau broker ou fournir des outils propriétaires.
⚠️ Rayon d'impact : N'exécutez jamais de commandes de sauvegarde ou de restauration sur un cluster de production sans runbook validé, test en staging isolé et chemin d'escalade d'astreinte.
Prérequis et Hypothèses
Versions de Kafka Couvertes
- 3.5.x : KRaft prêt pour la production (KIP-590), pas de stockage hiérarchisé en GA.
- 3.6.x : Stockage hiérarchisé en GA (KIP-405), instantanés de quorum de métadonnées améliorés (KIP-848).
- 3.7.x LTS : Support à long terme actuel ; inclut toutes les fonctionnalités 3.6 plus correctifs de stabilité.
Distinctions de Topologie
- KRaft : Le quorum de contrôleurs stocke les métadonnées dans la partition 0 de
__cluster_metadata; la sauvegarde nécessite un instantané de quorum ou la capture du journal du contrôleur. - ZooKeeper : L'ensemble stocke les métadonnées du cluster sous
/cluster/meta,/brokers,/config; sauvegarde via exportzkCli.sh.
Exigences d'Infrastructure
- Stockage bloc : AWS EBS, GCP Persistent Disk, Azure Managed Disk, ou PVC Kubernetes (Portworx, Longhorn, Rook-Ceph).
- Stockage objet : S3, GCS ou Azure Blob pour l'export de stockage hiérarchisé et l'archivage d'instantanés.
- Réseau : VPC peering dédié ou Transit Gateway pour la réplication inter-régions ; bande passante ≥ débit de production maximal (mesuré via
kafka.server.BrokerTopicMetrics.MessagesInPerSec).
Accès et Sécurité
- Accès CLI broker via scripts
kafka-*.shsur un hôte d'administration. - Jetons API admin ou propriétés client SASL/SSL (
$CLIENT_PROPS) avec ACLsDescribeConfigs,Read,ClusterActionpour le principal de sauvegarde. - Rôles IAM pour les API de snapshot cloud (
ec2:CreateSnapshot,s3:PutObject). - Lecture/écriture Schema Registry pour export/import des sujets.
- Chiffrement au repos (chiffrement EBS, classe de stockage PVC) et en transit (TLS 1.2+).
Planification de Capacité
- Staging local : 1,5×–2× la taille totale de
log.dirpour les exports de segments et le staging d'instantanés. - Rétention : Minimum 2× la fenêtre RPO pour les méthodes incrémentielles ; 7–30 jours pour les archives d'instantanés.
Architecture et Sélection de Stratégie
Choisissez une méthode en mappant vos RPO/RTO à la matrice de décision ci-dessous. Toutes les valeurs RPO/RTO supposent des clusters sains et une bande passante provisionnée ; validez dans votre environnement.
| Méthode | RPO | RTO | Complexité | Coût | Idéal Pour |
|---|---|---|---|---|---|
| MirrorMaker 2 (actif-actif) | < 1 min | < 5 min | Élevée | 2× cluster | Actif-actif multi-région, bascule géographique |
| MirrorMaker 2 (actif-passif) | < 1 min | < 15 min | Moyenne | 1,5× cluster | Reprise après sinistre avec standby chaud |
| Connect Replicator (Confluent) | < 5 min | < 15 min | Moyenne | Licence + 1,5× | Intégration Schema Registry, synchronisation ACL |
| Instantanés Volume/Bloc | < 1 h | < 4 h | Faible | Stockage snapshots | Récupération point-in-time, outillage minimal |
| Export Stockage Hiérarchisé (3.6+) | < 15 min | < 1 h | Moyenne | Stockage objet | Récupération sélective de sujets, archive conformité |
kafka-dump-log (manuel) | N/A | Heures | Faible | Disque local | Diagnostic corruption segment, récupération forensique |
Stratégie de Sauvegarde des Métadonnées
- KRaft (3.5+) :
kafka-metadata-quorum snapshotcapture l'état du quorum de contrôleurs ; exécutez sur le leader du contrôleur. - ZooKeeper : Exports
zkCli.sh get /cluster/meta+/brokers+/config; nécessite le quorum de l'ensemble.
Règles de Sélection des Sujets
- Toujours inclure :
__consumer_offsets,__transaction_state,_schemas(Schema Registry),__consumer_timestamps. - Inclure
__cluster_metadata(KRaft) uniquement via instantané de quorum, pas par copie de fichier. - Exclure les autres sujets internes
__*sauf pour une reprise après sinistre complète du cluster.
Garanties de Cohérence
- Cohérence au crash : Instantanés de volume sans quiesce ; le broker récupère depuis le dernier point de contrôle (risque : perte partielle de segments).
- Cohérence applicative : Mettez en pause les producteurs (
kafka-producer-perf-test --pause), attendez que le lag consommateur = 0, videz (producer.flush()), puis instantané. Ajoute 30–120s d'indisponibilité mais garantit l'intégrité des offsets/transactions.
Inventaire de Version et d'Environnement
Exécutez toutes les commandes de découverte en lecture seule ; capturez la sortie dans un fichier d'inventaire horodaté (inventory-$(date +%Y%m%d-%H%M%S).txt).
# Version et compatibilité de protocole
kafka-broker-api-versions --bootstrap-server $BOOTSTRAP --command-config $CLIENT_PROPS
# Santé du quorum KRaft (3.5+)
kafka-metadata-quorum --bootstrap-server $BOOTSTRAP --command-config $CLIENT_PROPS describe --status
# Configurations des sujets, facteur de réplication, cleanup.policy
kafka-topics --bootstrap-server $BOOTSTRAP --command-config $CLIENT_PROPS --describe
# Tailles des répertoires de logs, répertoires hors ligne
kafka-log-dirs --bootstrap-server $BOOTSTRAP --command-config $CLIENT_PROPS --describe --topic-list "$TOPICS"
# Métadonnées ZooKeeper (si applicable)
zookeeper-shell $ZK_HOST get /cluster/meta 2>/dev/null
zookeeper-shell $ZK_HOST ls /brokers/ids 2>/dev/null
Signaux de Vérification
- Le nombre de brokers depuis
kafka-broker-api-versionscorrespond à la taille attendue du cluster. - KRaft :
LeaderIdcorrespond à un votant ;CurrentVoters≥ 3 ;HighWatermarken progression. - ZooKeeper : Taille de l'ensemble impaire (≥3) ;
/cluster/metaretourne un JSON valide. - Aucun répertoire de logs hors ligne (
OfflineLogDirectoryCount = 0).
Procédures d'Exécution de Sauvegarde
Méthode A : MirrorMaker 2 (Reprise Après Sinistre Actif-Passif)
Prérequis
- Cluster cible accessible depuis les workers Connect MM2.
- Cluster Connect MM2 déployé (séparé des brokers source/cible).
offset-syncs.topic.replication.factor=3,heartbeats.topic.replication.factor=3,checkpoints.topic.replication.factor=3.- Clusters source et cible sur versions compatibles (MM2 depuis 2.4 ; 3.5+ recommandé).
Configuration du Connecteur (mm2-dr-connector.json)
{
"name": "mm2-dr-prod-to-dr",
"config": {
"connector.class": "org.apache.kafka.connect.mirror.MirrorSourceConnector",
"source.cluster.alias": "prod",
"target.cluster.alias": "dr",
"source.cluster.bootstrap.servers": "$PROD_BOOTSTRAP",
"target.cluster.bootstrap.servers": "$DR_BOOTSTRAP",
"source.cluster.security.protocol": "SASL_SSL",
"target.cluster.security.protocol": "SASL_SSL",
"source.cluster.sasl.mechanism": "SCRAM-SHA-512",
"target.cluster.sasl.mechanism": "SCRAM-SHA-512",
"topics": ".*",
"groups": ".*",
"emit.checkpoints.interval.seconds": "10",
"sync.topic.acls.enabled": "false",
"sync.group.offsets.enabled": "true",
"sync.group.offsets.interval.seconds": "60",
"offset-syncs.topic.replication.factor": "3",
"heartbeats.topic.replication.factor": "3",
"checkpoints.topic.replication.factor": "3",
"replication.factor": "3",
"tasks.max": "4"
}
}
Déploiement et Vérification
# Déploiement via Connect REST
curl -X POST -H "Content-Type: application/json" --data @mm2-dr-connector.json $CONNECT_URL/connectors
# Vérification du lag de synchronisation d'offsets = 0
kafka-consumer-groups --bootstrap-server $DR_BOOTSTRAP --command-config $DR_PROPS \
--group mm2-offset-syncs.prod --describe
✅ Signal de Vérification : Colonne LAG = 0 pour toutes les partitions ; CURRENT-OFFSET correspond au LOG-END-OFFSET source dans les 10s.
Méthode B : Kafka Connect Replicator (Confluent Platform)
Prérequis
- Licence Confluent Platform.
- Schema Registry sur les deux clusters avec mode de compatibilité de sujet compatible.
offset.topic.replication.factor=3,config.topic.replication.factor=3,status.topic.replication.factor=3.
Configuration du Connecteur (replicator-dr.json)
{
"name": "replicator-prod-to-dr",
"config": {
"connector.class": "io.confluent.connect.replicator.ReplicatorSourceConnector",
"topic.rename.format": "${topic}",
"src.kafka.bootstrap.servers": "$PROD_BOOTSTRAP",
"dest.kafka.bootstrap.servers": "$DR_BOOTSTRAP",
"src.kafka.security.protocol": "SASL_SSL",
"dest.kafka.security.protocol": "SASL_SSL",
"src.kafka.sasl.mechanism": "SCRAM-SHA-512",
"dest.kafka.sasl.mechanism": "SCRAM-SHA-512",
"src.consumer.group.id": "replicator-prod-to-dr",
"offset.topic.replication.factor": "3",
"config.topic.replication.factor": "3",
"status.topic.replication.factor": "3",
"topic.regex": ".*",
"tasks.max": "4"
}
}
Méthode C : Instantanés Volume/Bloc
Prérequis
- Broker arrêté OU quiesce du système de fichiers + pause d'écriture du contrôleur.
- Pour KRaft : mettez en pause les écritures de métadonnées du contrôleur avant l'instantané.
- Permissions IAM :
ec2:CreateSnapshot,ec2:CreateTags,ec2:DescribeVolumes.
Quiesce du Contrôleur KRaft (3.5+)
# Déclencher l'élection de leader préférée pour la partition 0 de __cluster_metadata pour vider les écritures en attente
kafka-leader-election --bootstrap-server $BOOTSTRAP --command-config $CLIENT_PROPS \
--election-type PREFERRED --topic __cluster_metadata --partition 0
# Attendre la fin de l'élection (surveiller ActiveControllerCount=1)
Instantané EBS avec Étiquettes de Cycle de Vie
aws ec2 create-snapshot \
--volume-id $VOL_ID \
--description "kafka-broker-$BROKER_ID-$(date +%s)" \
--tag-specifications 'ResourceType=snapshot,Tags=[{Key=KafkaBackup,Value=true},{Key=BrokerId,Value='$BROKER_ID'},{Key=Cluster,Value='$CLUSTER_NAME'},{Key=Environment,Value=production}]'
Procédure de Restauration
- Attachez le volume restauré à l'instance de remplacement au même chemin de périphérique (
/dev/xvdf→/var/lib/kafka/data). - Critique : Corrigez le décalage
broker.idavant le démarrage.
- ZooKeeper : Éditez
meta.propertiesdanslog.dir:broker.id=$ORIGINAL_BROKER_ID. - KRaft : Mettez à jour
controller.quorum.votersdansserver.propertiesou la configuration dynamique pour inclure l'ID du broker restauré.
- Démarrez le broker ; vérifiez
kafka-broker-api-versionset la récupération ISR (UnderReplicatedPartitions=0).
Méthode D : Export de Stockage Hiérarchisé (Kafka 3.6+)
Prérequis
remote.log.storage.system.enable=true.remote.storage.s3.bucket=$BUCKET(ou équivalent GCS/Azure).remote.log.metadata.manager.class.name=org.apache.kafka.server.log.remote.storage.RemoteLogMetadataManager.- Le broker a
s3:GetObject,s3:PutObjectsur le bucket.
Commande d'Export
Vérifiez la syntaxe dans la documentation de la version cible ; le CLI évolue.
kafka-tiered-storage --bootstrap-server $BOOTSTRAP --command-config $CLIENT_PROPS \
--export --topic $TOPIC --partition $PART \
--remote-path s3://$BUCKET/export/$TOPIC-$PART-$(date +%s)
Import et Validation
kafka-tiered-storage --bootstrap-server $TARGET_BOOTSTRAP --command-config $TARGET_PROPS \
--import --remote-path s3://$BUCKET/export/$TOPIC-$PART-<timestamp> \
--local-log-dir /var/lib/kafka/data
# Valider l'intégrité de l'index
kafka-dump-log --files /var/lib/kafka/data/$TOPIC-$PART/00000000000000000000.log --verify-index-only
✅ Signal de Vérification : kafka-dump-log --verify-index-only sort avec le code 0 sans erreurs "Corrupt index".
Méthode E : kafka-dump-log (Inspection de Segments et Récupération Point-in-Time)
Cas d'Usage : Diagnostic de corruption de segment, récupération forensique de segments supprimés, ligne de base de checksum.
# Dump d'itération profonde avec enregistrements de données (pour checksum)
kafka-dump-log --files /var/lib/kafka/data/$TOPIC-$PART/00000000000000000000.log \
--print-data-log --deep-iteration > /backup/$TOPIC-$PART-$(date +%s).dump
# Pipeline de checksum (supprimer les lignes d'en-tête)
kafka-dump-log --files /var/lib/kafka/data/$TOPIC-$PART/00000000000000000000.log \
--deep-iteration | grep -v "^Dumping" | sha256sum
Vérification et Diagnostics
Exécutez la validation après chaque opération de sauvegarde ou restauration. Automatisez dans CI/CD ou cron avec alertes sur échec.
Commandes de Validation Post-Sauvegarde
| Méthode | Commande de Validation | Critère de Succès |
|---|---|---|
| MM2 | kafka-consumer-groups --bootstrap-server $TARGET --group mm2-offset-syncs.$SOURCE --describe | LAG de toutes les partitions = 0 |
| Instantané Volume | kafka-log-dirs --bootstrap-server $RESTORED --describe --topic-list "$TOPICS" | Tailles des répliques correspondent à la source ±1% |
| Export Hiérarchisé | kafka-dump-log --files $RESTORED_LOG --verify-index-only | Code de sortie 0, aucune erreur |
kafka-dump-log | kafka-dump-log --files $LOG --deep-iteration | grep -v "^Dumping" | sha256sum | Checksum correspond à la ligne de base source |
Validation des Offsets de Groupe de Consommateurs
# Export des offsets source
kafka-consumer-groups --bootstrap-server $SRC --command-config $SRC_PROPS \
--group $CONSUMER_GROUP --export > /backup/offsets-$CONSUMER_GROUP-$(date +%s).json
# Import vers cible (dry-run d'abord avec --dry-run)
kafka-consumer-groups --bootstrap-server $TGT --command-config $TGT_PROPS \
--import --input-file /backup/offsets-$CONSUMER_GROUP.json --reset-offsets --execute
✅ Signal de Vérification : Le groupe de consommateurs cible montre un LAG stable (pas en croissance), aucun membre UNKNOWN_MEMBER_ID.
Intégrité du Journal de Transactions
# Compter les enregistrements de contrôle (marqueurs commit/abort) dans __transaction_state
kafka-dump-log --files /var/lib/kafka/data/__transaction_state-0/00000000000000000000.log \
--print-data-log --deep-iteration | grep -c "control record"
Comparez le nombre entre la source et la réplica restaurée ; un décalage indique une perte de transactions.
Alertes de Surveillance pour la Santé des Sauvegardes
- Chute de
kafka_server_BrokerTopicMetrics_MessagesInPerSec> 50% pendant la fenêtre de sauvegarde → investiguer la limitation. kafka_controller_KafkaController_ActiveControllerCount!= 1 → instabilité du contrôleur.kafka_log_LogManager_OfflineLogDirectoryCount> 0 → défaillance disque pendant l'instantané.- Personnalisé : MM2
mm2-source-task-lag-max> 10000 pendant 5min → blocage de réplication.
Modes de Défaillance et Récupération
Tâche Connecteur MM2 en ÉCHEC
Symptôme : Statut de tâche FAILED dans le sujet connect-statuses ; logs montrent TimeoutException ou AuthenticationException.
Diagnostic :
# Vérifier le statut de la tâche
curl $CONNECT_URL/connectors/mm2-dr-prod-to-dr/status | jq '.tasks[] | select(.state=="FAILED")'
Récupération :
# Redémarrer la tâche spécifique
curl -X POST $CONNECT_URL/connectors/mm2-dr-prod-to-dr/tasks/0/restart
# Si persistant, vérifier connectivité source/cible, ACLs, et réplication du sujet offset-syncs
Restauration d'Instantané : Décalage d'ID Broker
Symptôme : Le broker démarre mais ne peut pas rejoindre le cluster ; logs montrent "Broker ID X already registered" ou décalage de votant KRaft.
Récupération :
- ZooKeeper : Arrêtez le broker, éditez
meta.propertiesdanslog.dirvers l'broker.idoriginal, redémarrez. - KRaft : Mettez à jour
controller.quorum.votersdansserver.propertiessur tous les contrôleurs pour inclure l'ID du broker restauré ; redémarrage progressif des contrôleurs.
Réduction ISR Après Restauration
Symptôme : UnderReplicatedPartitions > 0 persiste > 5 min après redémarrage du broker.
Récupération :
# Déclencher l'élection de leader préférée pour toutes les partitions
kafka-leader-election --bootstrap-server $BOOTSTRAP --command-config $CLIENT_PROPS \
--election-type PREFERRED --all-topic-partitions
# Surveiller jusqu'à UnderReplicatedPartitions=0
watch -n 10 "kafka-server-metrics --bootstrap-server $BOOTSTRAP --command-config $CLIENT_PROPS | grep UnderReplicatedPartitions"
Tempête de Réinitialisation d'Offsets
Symptôme : Les consommateurs retraiter depuis earliest après restauration ; pic de lag, traitement en double.
Récupération :
# Mettre en pause les consommateurs (changement de config ou SIGSTOP)
# Réinitialiser les offsets à la position committée (pas earliest)
kafka-consumer-groups --bootstrap-server $TARGET --command-config $TARGET_PROPS \
--group $GROUP --reset-offsets --to-current --execute --dry-run
# Vérifier, puis exécuter sans --dry-run
# Reprendre les consommateurs
Incompatibilité Schema Registry
Symptôme : Les consommateurs échouent avec SchemaParseException ou IncompatibleSchemaException après bascule DR.
Récupération :
# Exporter tous les sujets depuis le SR source
sr-cli subjects | xargs -I{} sr-cli get-schema {} > /backup/sr-subjects-$(date +%s).json
# Importer vers le SR cible (vérifier le mode de compatibilité d'abord)
cat /backup/sr-subjects.json | jq -c '.[]' | while read subject; do
sr-cli register --subject "$(echo $subject | jq -r .subject)" --schema "$(echo $subject | jq -r .schema)"
done
Procédures de Retour Arrière Par Méthode
| Méthode | Déclencheur de Retour Arrière | Étapes de Retour Arrière |
|---|---|---|
| MM2 | Lag du cluster cible > seuil RTO | curl -X PUT $CONNECT_URL/connectors/mm2-dr-prod-to-dr/pause ; vider le lag via vidage consommateur ; reprendre quand lag < seuil |
| Instantané Volume | Restauration corrompue détectée | Détacher le volume restauré ; rattacher l'original ; redémarrer le broker |
| Export Hiérarchisé | Échec de validation d'import | Supprimer les segments importés de log.dir ; relancer l'import depuis l'export précédent |
kafka-dump-log | Décalage de checksum | Jeter le segment restauré ; recopier depuis la ligne de base source |
Liste de Contrôle Opérationnelle
- [ ] Inventaire capturé et version vérifié (commandes Section 3 exécutées, sortie archivée).
- [ ] Méthode de sauvegarde sélectionnée selon la matrice RPO/RTO (décision Section 2 documentée).
- [ ] Prérequis remplis : marge de stockage (1,5×–2×), bande passante réseau, rôles IAM, accès Schema Registry.
- [ ] Test de restauration exécuté en staging dans les 30 derniers jours (runbook documenté avec horodatages).
- [ ] Scripts de vérification automatisés : comparaison de checksum, lag consommateur, santé ISR, compte de journal de transactions.
- [ ] Runbook documenté avec rayon d'impact, étapes de retour arrière, contacts d'astreinte, plan de communication.
- [ ] Secrets gérés via Vault/Secrets Manager ; aucun texte en clair dans les configs ou l'historique des commandes.
- [ ] Alertes de surveillance ajustées : succès/échec du job de sauvegarde, lag de réplication, santé du contrôleur, répertoires de logs hors ligne.
- [ ] Instantané de quorum de métadonnées KRaft planifié quotidiennement (
kafka-metadata-quorum snapshot --output-file /backup/kraft-snapshot-$(date +%s).bin). - [ ] Export/import des sujets Schema Registry testé trimestriellement.
Conclusion
La sauvegarde de Kafka en production n'est pas un outil unique mais un portefeuille de stratégies selon la version, alignées sur des objectifs RPO/RTO explicites. MirrorMaker 2 offre un RPO quasi nul pour les topologies actif-actif et actif-passif mais exige une maturité opérationnelle ; les instantanés de volume apportent la simplicité avec un RPO à l'heure ; l'export de stockage hiérarchisé (3.6+) comble l'écart pour la récupération sélective. Les instantanés de quorum de métadonnées KRaft (KIP-848) et l'inclusion de __consumer_offsets / __transaction_state sont non négociables pour une restauration cohérente. Chaque méthode requiert une vérification automatisée — correspondance de checksum, lag=0, ISR stable, intégrité du journal de transactions — avant de déclarer une sauvegarde valide. Testez les restaurations en staging mensuellement, documentez les déclencheurs de retour arrière et gardez les secrets hors de l'historique des commandes. La seule sauvegarde qui compte est celle que vous avez restaurée avec succès.