Ce guide apprend aux ingénieurs plateforme à implémenter un pilote CI/CD étroit et sûr en production pour Apache Kafka qui et les ACLs de manière déclarative. Il fournit une base complète et exécutable : inventaire d'environnement, structure de dépôt, scripts d'application idempotents avec garde-fous, commandes de vérification avec sorties attendues, runbook des modes de défaillance, procédures de retour arrière et liste de contrôle opérationnelle. L'objectif est d'établir des patterns d'automatisation observables et réversibles que les équipes peuvent étendre aux quotas, schémas et environnements supplémentaires sans retravail.
Prérequis, versions et hypothèses
Version minimale des brokers Kafka : 3.6.x utilisée dans les exemples ; parité de version CLI obligatoire. Outils CLI requis : kafka-topics.sh, kafka-acls.sh, kafka-configs.sh, kafka-consumer-groups.sh, kafka-reassign-partitions.sh, kcat (optionnel mais recommandé). Hypothèses de sécurité : SASL_SSL avec SCRAM-SHA-512 (primaire), alternative mTLS notée. Exigences runner : accessibilité réseau, résolution DNS, synchronisation temporelle (Kerberos/SCRAM), CLI Kafka installée correspondant à la version broker. Dépôt Git avec branches protégées, revues requises, commits signés. Intégration gestionnaire de secrets (Vault, AWS Secrets Manager, GCP Secret Manager, GitHub Environments).
Architecture et flux de données
État déclaratif (Git) → Pipeline CI (plan → apply → verify) → Cluster Kafka. Séparation : environments/<env>/ isole le rayon d'impact. Scripts d'application idempotents comme unique chemin de mutation ; pas de changements CLI manuels. Portes de vérification : connectivité → apply → describe → smoke → lag → promotion.
Gestion de la configuration et des secrets
Le pattern env.sh est étendu avec KAFKA_SECURITY_PROTOCOL, KAFKA_JAAS_CONFIG (ou fichier --command-config), KAFKA_SSL_TRUSTSTORE_LOCATION, KAFKA_SSL_KEYSTORE_LOCATION, KAFKA_SASL_MECHANISM. Exemple client.properties pour --command-config (préféré aux variables d'env pour les secrets). CI : injection des secrets à l'exécution via échange de jeton OIDC ou variables d'environnement masquées ; jamais logger les secrets.
# client.properties template (SASL_SSL)
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="${KAFKA_USER}" password="${KAFKA_PASSWORD}";
ssl.truststore.location=/opt/kafka/secrets/truststore.jks
ssl.truststore.password=${TRUSTSTORE_PASSWORD}
ssl.keystore.location=/opt/kafka/secrets/keystore.jks
ssl.keystore.password=${KEYSTORE_PASSWORD}
ssl.key.password=${KEY_PASSWORD}
# client.properties template (mTLS)
security.protocol=SSL
ssl.truststore.location=/opt/kafka/secrets/truststore.jks
ssl.truststore.password=${TRUSTSTORE_PASSWORD}
ssl.keystore.location=/opt/kafka/secrets/keystore.jks
ssl.keystore.password=${KEYSTORE_PASSWORD}
ssl.key.password=${KEY_PASSWORD}
Scripts d'application durcis
safe-apply-topics.sh : ajoute un drapeau --dry-run qui imprime les actions planifiées (création, augmentation partitions, changements config) sans exécuter ; valide RF ≤ nombre de brokers ; rejette les clés config en lecture seule (ex. cleanup.policy sur topics compactés) ; utilise kafka-configs.sh --describe pour différencier avant alter.
safe-apply-acls.sh : supporte resource-pattern-type (Literal|Prefixed) depuis le fichier ; supporte liste --operation (Read,Write,Create,Delete,Alter,Describe,ClusterAction,IdempotentWrite,All) ; ajoute --dry-run listant le diff kafka-acls.sh --list ; gère format principal User: vs Group:.
Les deux scripts : émettent des logs JSON structurés (timestamp, level, env, resource, action, status, correlation_id) pour l'observabilité.
#!/usr/bin/env bash
# safe-apply-topics.sh (extrait durci)
set -euo pipefail
source "$(dirname "$0")/env.sh"
DRY_RUN=false
"${1:-}" == "--dry-run" && { DRY_RUN=true; shift; }
env_dir=${1:?usage: $0 [--dry-run] <environments/dev|environments/prod>}
CORR_ID=$(uuidgen)
log_json() { jq -n --arg ts "$(date -u +%Y-%m-%dT%H:%M:%SZ)" --arg lvl "$1" --arg env "$2" --arg res "$3" --arg act "$4" --arg sts "$5" --arg cid "$CORR_ID" '{timestamp:$ts,level:$lvl,env:$env,resource:$res,action:$act,status:$sts,correlation_id:$cid}'; }
create_or_update_topic() {
local f="$1" name partitions rf
name=$(grep '^name=' "$f" | sed 's/^name=//')
partitions=$(grep '^partitions=' "$f" | sed 's/^partitions=//')
rf=$(grep '^replication.factor=' "$f" | sed 's/^replication.factor=//')
mapfile -t configs < <(grep '^config\.' "$f" | sed 's/^config\.//')
if kafka-topics.sh --bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --topic "$name" --describe >/dev/null 2>&1; then
echo "$(log_json INFO "$env_dir" "$name" "describe" "exists")"
current_p=$(kafka-topics.sh --bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --topic "$name" --describe | awk -F ':' '/PartitionCount/ {print $3}' | tr -d ' ')
if -n "$current_p" && "$partitions" -lt "$current_p" ; then
echo "$(log_json ERROR "$env_dir" "$name" "partition_decrease" "blocked")" >&2
return 1
fi
if "$partitions" -gt "$current_p" ; then
echo "$(log_json INFO "$env_dir" "$name" "partition_increase" "planned")"
$DRY_RUN || kafka-topics.sh --bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --alter --topic "$name" --partitions "$partitions"
fi
for c in "${configs[@]}"; do
key=${c%%=*}; val=${c#*=}
echo "$(log_json INFO "$env_dir" "$name" "config_alter" "planned")"
$DRY_RUN || kafka-configs.sh --bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --alter --entity-type topics --entity-name "$name" --add-config "$key=$val"
done
else
echo "$(log_json INFO "$env_dir" "$name" "create" "planned")"
$DRY_RUN || kafka-topics.sh --bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --create --topic "$name" --partitions "$partitions" --replication-factor "$rf" $(printf -- '--config %s ' "${configs[@]}")
fi
}
find "$env_dir/topics" -type f -name "*.properties" | while read -r f; do create_or_update_topic "$f"; done
Sortie --dry-run exemple (lignes JSON) :
{"timestamp":"2024-01-15T10:30:00Z","level":"INFO","env":"environments/dev","resource":"events.test.v1","action":"create","status":"planned","correlation_id":"a1b2-c3d4"}
{"timestamp":"2024-01-15T10:30:01Z","level":"INFO","env":"environments/dev","resource":"events.test.v1","action":"config_alter","status":"planned","correlation_id":"a1b2-c3d4"}
#!/usr/bin/env bash
# safe-apply-acls.sh (extrait supportant Prefixed)
set -euo pipefail
source "$(dirname "$0")/env.sh"
DRY_RUN=false
"${1:-}" == "--dry-run" && { DRY_RUN=true; shift; }
env_dir=${1:?usage: $0 [--dry-run] <environments/dev|environments/prod>}
apply_acl_file() {
local f="$1" principal operation resource pattern_type host group ops
principal=$(grep '^principal=' "$f" | sed 's/^principal=//')
operation=$(grep '^operation=' "$f" | sed 's/^operation=//')
resource=$(grep '^resource=' "$f" | sed 's/^resource=//')
pattern_type=$(grep '^resource.pattern.type=' "$f" | sed 's/^resource.pattern.type=//' || echo "Literal")
host=$(grep '^host=' "$f" | sed 's/^host=//' || echo "*")
group=$(grep '^consumer.group=' "$f" | sed 's/^consumer.group=//' || true)
IFS=',' read -ra ops <<< "$operation"
IFS=':' read -r rtype rname <<< "$resource"
for op in "${ops[@]}"; do
args=(--bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --add --allow-principal "$principal" --operation "$op" --resource-pattern-type "$pattern_type")
"$rtype" == "Topic" && args+=(--topic "$rname")
"$rtype" == "Group" && args+=(--group "$rname")
-n "$host" && args+=(--host "$host")
$DRY_RUN || kafka-acls.sh "${args[@]}"
done
if [[ -n "$group" && "$rtype" == "Topic" && " ${ops[*]} " =~ " Read " ]]; then
args=(--bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --add --allow-principal "$principal" --operation Read --group "$group" --resource-pattern-type Literal)
$DRY_RUN || kafka-acls.sh "${args[@]}"
fi
}
find "$env_dir/acls" -type f -name "*.acl" | while read -r f; do apply_acl_file "$f"; done
Vérification et critères d'acceptation
Script automatisé verify-apply.sh <env> exécutant tous les contrôles et sortant en erreur si écart. Seuils : lag consommateur < 10 pour groupe canari avec timeout 30s.
#!/usr/bin/env bash
# verify-apply.sh
set -euo pipefail
source "$(dirname "$0")/env.sh"
env_dir=${1:?usage: $0 <environments/dev|environments/prod>}
topic=$(grep '^name=' "$env_dir/topics"/*.properties | head -1 | sed 's/.*=//')
group=$(grep 'consumer.group=' "$env_dir/acls"/*.acl | head -1 | sed 's/.*=//' || echo "grp.canary")
echo "--- Vérification topic ---"
kafka-topics.sh --bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --describe --topic "$topic"
echo "--- Vérification ACLs ---"
kafka-acls.sh --bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --list --topic "$topic"
echo "--- Smoke test ---"
scripts/smoke-produce-consume.sh "$topic"
echo "--- Lag consommateur ---"
timeout 30s bash -c "
while true; do
lag=$(kafka-consumer-groups.sh --bootstrap-server \"$KAFKA_BOOTSTRAP\" --command-config \"$KAFKA_COMMAND_CONFIG\" --describe --group \"$group\" 2>/dev/null | awk '/^$topic/ {sum+=\$NF} END {print sum+0}')
\$lag -lt 10 && { echo \"LAG OK: \$lag\"; exit 0; }
sleep 2
done
" || { echo "LAG threshold exceeded or timeout"; exit 1; }
Sorties attendues :
kafka-topics.sh --describe:PartitionCount: 12,ReplicationFactor: 3,Config: retention.ms=604800000,cleanup.policy=deletekafka-acls.sh --list:Principal: User:svc-producer, Operation: Write, Resource: Topic:events.test.v1, Pattern: LITERALetPrincipal: User:svc-consumer, Operation: Read, Resource: Topic:events.test.v1, Pattern: LITERAL+Group: grp-orders-downstream- Smoke test :
[OK] Smoke test passed for events.test.v1 - Lag :
LAG OK: 0
Modes de défaillance, dépannage et runbooks
- Symptôme — Cause probable — Action immédiate — Étape de reprise
- Connexion refusée / timeout — Bootstrap ou firewall erronés — Vérifier
KAFKA_BOOTSTRAPet réseau — Corriger inventaire et relancer check - Authorization failed — ACLs/identifiants manquants — Vérifier principal/ACLs via
--list— Appliquer ACLs correctes, relancer smoke test - Échec création topic (RF) — RF > nb brokers — Ajuster RF dans
.properties— Recréer avec RF valide ou ajouter brokers - Demande baisse partitions — Mutation dangereuse non supportée — Blocage par design (script refuse) — Augmenter seulement ; sinon nouveau topic + migration
- Topic config drift détecté — Changement manuel hors CI —
kafka-configs.sh --describepour identifier — Ré-appliquer depuis Git ou promouvoir changement manuel vers Git - Augmentation partitions bloquée — Réplicas non convergents (ISR mismatch) —
kafka-topics.sh --describemontre Replicas vs ISR — Attendre sync ISR ou déclencher réassignation - Délai propagation ACL — Cache authorizer (défaut 30s) — Attendre
authorizer.cache.seconds— Redémarrer brokers si SimpleAclAuthorizer sans expiry cache - Schema registry compatibility failure — Schéma incompatible (extension) — Bloquer déploiement, lancer check compatibilité CI — Corriger schéma, valider compatibilité AVRO/PROTOBUF
Procédures de retour arrière (détail opérationnel)
- Rollback config :
git revert <sha> && ./safe-apply-topics.sh env(idempotent). - Remplacement topic : créer
topic.v2, dual-write via config producteur, miroir avec MirrorMaker2 ou replicateur custom, basculer consommateurs, supprimertopic.v1après TTL. - Correction RF : générer JSON réassignation avec
kafka-reassign-partitions.sh --generate, réviser, exécuter--execute, vérifier--verify. - ACL over-grant : appliquer ACL restrictive, puis
kafka-acls.sh --removepour la large ; vérifier avec--list. - Rollback environnement complet : tag Git
pre-change-<timestamp>, revert au tag, re-appliquer.
Exemple JSON réassignation pour augmentation RF :
{
"version": 1,
"partitions": [
{"topic": "events.test.v1", "partition": 0, "replicas": [1,2,3,4]},
{"topic": "events.test.v1", "partition": 1, "replicas": [2,3,4,1]}
]
}
Commandes : kafka-reassign-partitions.sh --bootstrap-server $KAFKA_BOOTSTRAP --command-config $KAFKA_COMMAND_CONFIG --reassignment-json-file reassign.json --execute puis --verify.
Vérification lag avec seuil et timeout :
kafka-consumer-groups.sh --bootstrap-server "$KAFKA_BOOTSTRAP" --command-config "$KAFKA_COMMAND_CONFIG" --describe --group grp-orders-downstream | awk '/events.test.v1/ {print $NF}'
Pipeline CI/CD (grade production)
GitHub Actions avec jobs plan/apply/verify/rollback, environnements, secrets, approbations.
name: kafka-cicd
on:
push:
paths:
- 'environments/**'
workflow_dispatch:
inputs:
environment:
type: choice
options: [dev, prod]
action:
type: choice
options: [plan, apply, rollback]
jobs:
plan:
runs-on: self-hosted
environment: ${{ github.event.inputs.environment || 'dev' }}
timeout-minutes: 10
steps:
- uses: actions/checkout@v4
- name: Setup Kafka CLI
uses: docker://apache/kafka:3.6.1
with:
entrypoint: ''
- name: Plan topics
run: scripts/safe-apply-topics.sh --dry-run environments/${{ github.event.inputs.environment || 'dev' }}
- name: Plan ACLs
run: scripts/safe-apply-acls.sh --dry-run environments/${{ github.event.inputs.environment || 'dev' }}
- name: Upload plan artifact
uses: actions/upload-artifact@v4
with:
name: kafka-plan
path: plan-output/
apply:
needs: plan
if: github.event.inputs.action == 'apply' || github.ref == 'refs/heads/main'
runs-on: self-hosted
environment: prod
timeout-minutes: 15
steps:
- uses: actions/checkout@v4
- name: Apply topics
env:
KAFKA_BOOTSTRAP: ${{ secrets.KAFKA_BOOTSTRAP_PROD }}
KAFKA_COMMAND_CONFIG: ${{ secrets.KAFKA_COMMAND_CONFIG_PROD }}
run: scripts/safe-apply-topics.sh environments/prod
- name: Apply ACLs
env:
KAFKA_BOOTSTRAP: ${{ secrets.KAFKA_BOOTSTRAP_PROD }}
KAFKA_COMMAND_CONFIG: ${{ secrets.KAFKA_COMMAND_CONFIG_PROD }}
run: scripts/safe-apply-acls.sh environments/prod
verify:
needs: apply
runs-on: self-hosted
environment: prod
steps:
- uses: actions/checkout@v4
- name: Verify
env:
KAFKA_BOOTSTRAP: ${{ secrets.KAFKA_BOOTSTRAP_PROD }}
KAFKA_COMMAND_CONFIG: ${{ secrets.KAFKA_COMMAND_CONFIG_PROD }}
run: scripts/verify-apply.sh environments/prod
rollback:
if: github.event.inputs.action == 'rollback'
runs-on: self-hosted
environment: prod
steps:
- uses: actions/checkout@v4
- name: Rollback to previous tag
run: |
git fetch --tags
PREV_TAG=$(git tag -l 'pre-change-*' --sort=-creatordate | head -1)
git checkout $PREV_TAG
scripts/safe-apply-topics.sh environments/prod
scripts/safe-apply-acls.sh environments/prod
Runner setup : runners auto-hébergés avec version CLI Kafka épinglée (image Docker apache/kafka:3.6.1 ou custom). Réseau : runners dans même VPC/sous-réseau que brokers, groupes sécurité autorisant port SASL_SSL.
Durcissement sécurité
Comptes de service moindre privilège : principaux séparés pour CI (admin), producteurs, consommateurs. ACLs pour principal CI : Cluster:Alter,Describe,ClusterAction + Topic:Create,Alter,Describe,Read,Write sur topics gérés uniquement (pattern Prefixed). Audit logging : activer logging authorizer Kafka, expédier vers SIEM. Rotation identifiants : automatiser rotation SCRAM via kafka-configs.sh --alter --entity-type users.
Performance et considérations d'échelle
Application par lots topics/ACLs : scripts traitent fichiers séquentiellement ; pour centaines de ressources, paralléliser avec xargs -P et rate-limiter. Éviter effet de meute : échelonner augmentations partitions across topics. Surveiller métriques controller pendant applies : kafka.controller:type=KafkaController,name=ActiveControllerCount, OfflinePartitionsCount.
Observabilité et audit
ID de corrélation passé CI → scripts → Kafka client.id. Logs structurés expédiés vers Loki/Elastic ; dashboards pour durée apply, taux succès, détection drift. SHA commit Git tagué sur ressources Kafka via kafka-configs.sh --alter --add-config "git.commit=<sha>" (informatif).
Exemple ligne log structuré (JSON) :
{"timestamp":"2024-01-15T10:30:05Z","level":"INFO","env":"environments/prod","resource":"events.test.v1","action":"partition_increase","status":"success","correlation_id":"a1b2-c3d4","git_commit":"f4e3d2a1"}
Scénario technique réaliste
L'équipe ajoute un nouveau flux d'événements orders.enriched.v1 avec 12 partitions, RF=3, rétention 7 jours, producteur svc-orders-enricher, groupe consommateur grp-orders-downstream.
Parcours : créer fichiers dans environments/dev/topics/orders.enriched.v1.properties et acls/producer.orders.enriched.v1.acl + consumer.orders.enriched.v1.acl, ouvrir PR, CI exécute plan (affiche création + 2 ACLs), réviseur approuve, merge déclenche apply sur dev, vérification passe, promotion vers prod via dispatch manuel.
Fichier topic exemple :
name=orders.enriched.v1
partitions=12
replication.factor=3
config.cleanup.policy=delete
config.retention.ms=604800000
config.min.insync.replicas=2
Fichier ACL producteur :
principal=User:svc-orders-enricher
operation=Write,Describe
resource=Topic:orders.enriched.v1
resource.pattern.type=Literal
host=*
Fichier ACL consommateur :
principal=User:svc-orders-downstream
operation=Read,Describe
resource=Topic:orders.enriched.v1
resource.pattern.type=Literal
host=*
consumer.group=grp-orders-downstream
Liste de contrôle des opérations
- Étape — Action
- 1 — Revoir inventaire et identifiants
- 2 — Valider connectivité aux brokers
- 3 — Confirmer que seuls fichiers env cible ont changé
- 4 — Exécuter safe-apply-topics.sh pour env cible
- 5 — Exécuter safe-apply-acls.sh pour env cible
- 6 — Décrire topics et lister ACLs
- 7 — Lancer smoke-produce-consume.sh sur topic test
- 8 — Vérifier lag groupe canari
- 9 — En cas d'échec, rollback via revert + reapply
- 10 — Archiver commandes et sorties pour audit
Étendre au-delà du pilote
Quand le pilote est stable, étendre les mêmes patterns :
- Plus de topics et principals : fichiers petits et explicites, groupés par service.
- Quotas : gérer quotas par utilisateur avec
kafka-configs.shcôté clients. - Schémas sujets : piloter évolution via registre, même discipline « définir-relire-déployer ».
- Plateformes traitement et data streaming : connecter Apache NiFi, Apache Spark, sinks HDFS, ou jobs déclenchés par Apache Airflow. Pour chaque intégration, commencer par flux étroit et testable, puis élargir après vérification.
Conclusion
Automatiser Kafka avec CI/CD devient plus sûr en démarrant par un pilote étroit et mesurable puis en n'étendant qu'après avoir établi des voies de vérification et de reprise claires. Définissez l'état désiré de façon déclarative, appliquez avec des scripts idempotents qui refusent les opérations dangereuses, et vérifiez avec des commandes et sorties concrètes. En cas d'échec, utilisez les logs et le runbook : rollback de configs, remplacement de topic pour les changements incompatibles, et revérifications systématiques. Standardisez la structure de votre dépôt, encodez l'inventaire d'environnement, activez le pilote sur un premier environnement, puis étendez progressivement aux quotas, schémas et autres environnements. Ces pratiques rendent votre déploiement Kafka, votre pipeline Kafka et votre retour arrière Kafka plus sûrs, prévisibles et rapides à exécuter.