Aller au contenu
Riadh Mnasri
← Retour au blog
8 min de lecture

Kafka : diagnostiquer un consumer lag, méthodiquement

Kafka : diagnostiquer un consumer lag, méthodiquement

Une alerte tombe : « consumer lag élevé sur le groupe facturation ». Les notifications partent avec vingt minutes de retard. Le réflexe le plus répandu est d'augmenter le nombre de réplicas du consommateur. Il marche dans un cas précis, ne change rien dans plusieurs autres, et aggrave la situation dans au moins un. Cet article propose une démarche pour diagnostiquer un lag avant de toucher à quoi que ce soit.

Ce que mesure le lag#

Pour chaque partition, Kafka connaît deux positions : le dernier offset écrit par les producteurs (le log end offset) et le dernier offset commité par le groupe de consommateurs. Le lag est la différence entre les deux.

partition 0 :  [0][1][2][3][4][5][6][7][8][9]
                            ▲              ▲
                     offset commité    log end offset
                            └── lag = 5 ───┘

Deux remarques qui évitent de mauvaises conclusions.

Le lag se lit par partition, pas en total. Un lag total de 50 000 peut être réparti équitablement sur 12 partitions, ou concentré à 49 000 sur une seule. Ce sont deux problèmes différents.

Un lag en nombre de messages ne dit pas grand-chose seul. 10 000 messages de retard sur un topic qui en reçoit 100 000 par seconde, c'est un dixième de seconde. Sur un topic qui en reçoit 10 par minute, c'est deux semaines. Ce qui compte pour l'utilisateur, c'est le retard en temps, et surtout sa tendance : un lag stable, même élevé, veut dire que le consommateur suit le rythme. Un lag qui croît en continu veut dire qu'il consomme moins vite que ce qui arrive, et qu'il ne rattrapera jamais sans changement.

Première étape : regarder le groupe#

La commande livrée avec Kafka donne presque tout ce qu'il faut :

bash
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
  --describe --group facturation
TOPIC     PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG     CONSUMER-ID       HOST
factures  0          184220          184231          11      consumer-1-a3f…   /10.0.1.12
factures  1          179804          179810          6       consumer-1-a3f…   /10.0.1.12
factures  2          96115           143902          47787   consumer-2-7c1…   /10.0.1.17
factures  3          181377          181390          13      consumer-3-e08…   /10.0.1.21

Cette sortie oriente déjà le diagnostic. Ici, une seule partition est en retard, et les autres suivent. Ce n'est pas un problème de capacité globale.

Trois formes reviennent presque toujours :

Ce qu'on voitPiste principale
Toutes les partitions en retard, lag qui croît partoutLe traitement est globalement trop lent pour le débit
Une ou deux partitions en retard, les autres à jourClé chaude ou message problématique sur ces partitions
Colonne CONSUMER-ID vide, ou qui change d'une exécution à l'autreLe groupe est instable : rebalances en boucle

Cas 1 : le traitement est trop lent#

Le consommateur passe son temps à traiter, et le débit entrant dépasse ce qu'il peut absorber. Avant d'ajouter des réplicas, il faut savoir où passe le temps. Dans la grande majorité des cas, ce n'est pas Kafka, c'est ce que le consommateur fait de chaque message : une requête SQL par message, un appel HTTP synchrone, une sérialisation coûteuse.

Les leviers, dans l'ordre où je les essaie :

  1. Mesurer le temps de traitement par message. Un simple timer autour du traitement suffit. Si un message prend 40 ms et qu'il en arrive 50 par seconde et par partition, le calcul est vite fait : le consommateur ne peut pas suivre.
  2. Traiter par lots. max.poll.records (500 par défaut) détermine combien de messages arrivent à chaque poll. Un insert en base de 500 lignes coûte beaucoup moins cher que 500 inserts d'une ligne.
  3. Ajouter des consommateurs, mais seulement jusqu'au nombre de partitions. Dans un groupe, une partition est lue par un seul consommateur à la fois. Avec 4 partitions, un cinquième consommateur reste inactif.
  4. Augmenter le nombre de partitions, si on a atteint la limite du point 3. C'est une opération à planifier : pour un topic avec des clés, ajouter des partitions change la partition cible de chaque clé, et donc l'ordre de traitement pendant la transition.
Attention

Paralléliser à l'intérieur d'un consommateur (un pool de threads qui traite les messages d'un même poll en parallèle) augmente le débit, mais casse l'ordre par partition et complique la gestion des offsets : on ne peut commiter un offset que lorsque tous les messages précédents sont traités. Des bibliothèques comme Confluent Parallel Consumer le font correctement. Le coder à la main est une source classique de pertes de messages.

Cas 2 : une partition en retard#

Quand une seule partition traîne, ajouter des consommateurs ne sert à rien : cette partition reste lue par un seul d'entre eux. Deux causes dominent.

Une clé chaude. Kafka envoie tous les messages d'une même clé dans la même partition, pour garantir leur ordre. Si 40 % du trafic concerne un seul client (un gros compte, un identifiant par défaut comme "unknown", ou une clé nulle mal gérée), sa partition reçoit 40 % du trafic. On le vérifie en comptant les messages par clé sur un échantillon. La correction touche au modèle : choisir une clé plus fine si l'ordre n'est nécessaire qu'à un niveau plus fin, ou isoler le gros client dans un topic dédié.

Un message qui bloque. Un message mal formé, trop gros, ou qui déclenche une erreur que le consommateur réessaie indéfiniment. Le lag de la partition croît pendant que CURRENT-OFFSET ne bouge plus du tout. C'est le signe le plus clair : l'offset figé. La réponse est une politique d'erreur explicite : un nombre limité de tentatives, puis l'envoi du message dans un topic d'erreurs (dead letter topic) pour analyse, et le consommateur continue. Spring Kafka le propose avec DefaultErrorHandler et DeadLetterPublishingRecoverer.

Cas 3 : les rebalances en boucle#

C'est le cas où ajouter des consommateurs aggrave les choses. Un rebalance redistribue les partitions entre les membres du groupe. Pendant qu'il a lieu, la consommation s'arrête (complètement avec le protocole historique, partiellement avec le protocole coopératif). S'il se produit toutes les quelques minutes, le groupe passe plus de temps à se réorganiser qu'à consommer.

La cause la plus fréquente est un traitement plus long que max.poll.interval.ms (5 minutes par défaut). Le consommateur n'appelle pas poll à temps, le coordinateur le considère comme mort et l'exclut du groupe. Ses partitions sont réattribuées, les messages non commités sont relus par un autre consommateur, qui prend autant de temps, et le cycle recommence.

poll() → 500 messages × 800 ms = 400 s de traitement
       → dépasse max.poll.interval.ms (300 s)
       → le consommateur est exclu, rebalance
       → les 500 messages sont relus ailleurs → même durée → nouveau rebalance

Les signes : des logs contenant Member ... leaving group ou Attempt to heartbeat failed since group is rebalancing, et une colonne CONSUMER-ID qui change d'une exécution de la commande à l'autre.

Les corrections :

  • réduire max.poll.records pour que chaque lot tienne largement dans le délai ; c'est le correctif le plus rapide ;
  • accélérer le traitement (cas 1) ;
  • augmenter max.poll.interval.ms, en dernier recours, car un consommateur réellement bloqué mettra plus longtemps à être détecté ;
  • passer au rebalance coopératif (CooperativeStickyAssignor) pour que seules les partitions déplacées s'arrêtent, et utiliser l'appartenance statique (group.instance.id) pour qu'un redémarrage de pod pendant un déploiement ne déclenche pas de rebalance s'il revient à temps.

Les pauses de GC longues produisent le même symptôme, par un autre chemin : une JVM figée plusieurs secondes n'envoie plus de heartbeats. Si les rebalances coïncident avec des pics mémoire, c'est vers la JVM qu'il faut regarder.

Surveiller avant l'alerte#

La commande kafka-consumer-groups.sh sert au diagnostic, pas à la surveillance. En continu, deux sources :

  • côté consommateur, la métrique JMX records-lag-max, exposée automatiquement par le client Java et reprise par Micrometer dans une application Spring Boot ;
  • côté cluster, un exporter Prometheus qui calcule le lag de tous les groupes sans dépendre des consommateurs, ce qui reste utile quand ils sont tous arrêtés.

Pour l'alerte, je préfère une règle sur la tendance plutôt qu'un seuil fixe : « le lag de ce groupe augmente depuis 15 minutes » plutôt que « le lag dépasse 10 000 ». Le seuil fixe réveille la nuit pour un pic que le consommateur aurait rattrapé seul en deux minutes, et reste silencieux sur un topic à faible débit où 500 messages représentent une journée de retard.

La démarche en résumé#

Le lag monte
 ├─ sur toutes les partitions ?
 │    └─ temps de traitement par message × débit > capacité → lots, puis réplicas (≤ partitions)
 ├─ sur une partition ?
 │    ├─ offset figé → message bloquant → retry limité + dead letter topic
 │    └─ offset qui avance lentement → clé chaude → revoir la clé
 └─ consommateurs qui changent / logs de rebalance ?
      └─ traitement > max.poll.interval.ms → réduire max.poll.records, coopératif, statique

Un lag n'est pas un problème de Kafka, c'est un symptôme. Kafka fait exactement ce qu'on lui demande : il garde les messages en attendant qu'on les lise. La cause se trouve presque toujours dans le consommateur ou dans le choix des clés, et la commande --describe, lue partition par partition, suffit à savoir dans quelle branche de l'arbre on se trouve.