Cloud

Azure Event Hubs : diagnostiquer le retard consommateur avant de scaler ou rejouer

Un runbook de production pour séparer pression d’ingestion, throttling, partition chaude, panne du consumer et dérive de checkpoint avant d’augmenter la capacité ou de rejouer des événements.

13 août 2026 azureevent-hubsstreamingconsumer-lagcheckpointpartitionsazure-monitorkqlobservabilityautomationrunbookrollbackproduction

Un consumer Azure Event Hubs prend du retard. Les événements arrivent encore, mais les traitements métier apparaissent plusieurs minutes plus tard et l’écart continue de grandir. Augmenter les Throughput Units, ajouter des instances ou remettre le checkpoint en arrière peut sembler prudent. Chacune de ces actions peut pourtant aggraver l’incident : davantage de consumers ne crée pas de parallélisme au-delà du nombre de partitions utiles, plus de capacité Event Hubs ne répare pas une dépendance aval lente, et un replay non borné peut dupliquer des effets déjà validés.

Le cas fil rouge est un pipeline qui alimente un système de détection en quasi temps réel. Un consumer group de production traite les événements par partition, écrit un checkpoint après succès et appelle une base aval. L’objectif du runbook est de décider à partir de preuves : corriger le consumer, lever un throttling de namespace, rééquilibrer la charge, effectuer un replay borné ou rollbacker le dernier changement.

Figer le symptôme et le périmètre

Commencez par nommer le flux exact. Event Hubs partage la capacité au niveau du namespace, tandis que l’ownership et les checkpoints se lisent par event hub, consumer group et partition. Une vue agrégée peut donc masquer un autre producteur qui consomme la capacité ou une seule partition qui accumule tout le retard.

yaml event-hubs-incident-scope.yml
incident: inc-20260813-014
namespace: <event-hubs-namespace>
event_hub: <event-hub-name>
consumer_group: <consumer-group>
environment: production
first_visible_delay_utc: <timestamp>
expected_end_to_end_delay_seconds: <slo>
observed_delay_seconds: <value>

consumer:
application: <service-name>
version: <release-id>
instances: <count>
checkpoint_store: <storage-account-or-provider>

preserve:
- incoming and outgoing message rates
- throttled requests and service errors
- lag or sequence gap per partition
- checkpoint age per partition
- consumer errors, duration and retries
- downstream latency and saturation
- deployment and configuration changes

Gelez les déploiements, le scale automatique applicatif et les opérations manuelles de checkpoint pendant la collecte initiale. Conservez aussi la rétention restante. Si le plus ancien événement non traité approche de l’expiration, la fenêtre de décision devient une donnée d’incident à part entière.

Construire une matrice par partition

Le retard n’est pas une valeur unique. Pour chaque partition, comparez la dernière séquence publiée, la dernière séquence traitée avec succès, le checkpoint persistant, l’owner actuel et l’âge du dernier progrès. Un consumer qui lit vite mais checkpoint rarement n’a pas le même problème qu’un consumer qui ne reçoit plus rien.

text partition-evidence-matrix.txt
Pour chaque partition
Event Hubs
  derniere sequence ou offset disponible
  debit entrant et sortant observe
  erreurs ou throttling dans la fenetre

Consumer
  owner actuel et derniere prise de lease
  derniere sequence recue
  derniere sequence traitee avec succes
  duree p50 p95 p99 et erreurs de traitement

Checkpoint
  sequence ou offset persistant
  heure de derniere mise a jour
  ecart avec le dernier succes applicatif

Aval
  latence, erreurs, saturation et limites
  identifiant idempotent ou effet deja applique

Classification
  retard uniforme sur toutes les partitions
  une partition chaude
  lecture active mais traitement lent
  traitement reussi mais checkpoint bloque
  ownership instable ou partitions non attribuees

En Premium et Dedicated, les application metrics logs peuvent exposer notamment ConsumerLag et OffsetCommit. Dans les autres configurations, ou lorsque cette télémétrie n’est pas activée, instrumentez le consumer : séquence reçue, séquence validée, checkpoint écrit et temps de traitement. Ne déduisez pas un lag précis uniquement de IncomingMessages - OutgoingMessages : plusieurs consumer groups peuvent lire le même flux et les compteurs sont agrégés.

Prouver ou exclure une limite de capacité

Regardez le namespace avant l’application : IncomingMessages, OutgoingMessages, IncomingBytes, OutgoingBytes, ThrottledRequests, ServerErrors et UserErrors. Comparez la fenêtre d’incident à une période saine et conservez la granularité la plus fine disponible. Les Throughput Units du niveau Standard sont partagées par tous les event hubs du namespace ; une charge voisine peut donc être causale.

bash 01-event-hubs-platform-metrics.sh
set -euo pipefail

RESOURCE_ID="/subscriptions/<subscription-id>/resourceGroups/<resource-group>/providers/Microsoft.EventHub/namespaces/<namespace>"
START_UTC="<incident-start-utc>"
END_UTC="<incident-end-utc>"

for METRIC in IncomingMessages OutgoingMessages IncomingBytes OutgoingBytes ThrottledRequests ServerErrors UserErrors; do
az monitor metrics list   --resource "$RESOURCE_ID"   --metric "$METRIC"   --interval PT1M   --start-time "$START_UTC"   --end-time "$END_UTC"   --aggregation Total   --output json > "metric-$METRIC.json"
done

az eventhubs namespace show --ids "$RESOURCE_ID" --query '{sku:sku.name,capacity:sku.capacity,autoInflate:isAutoInflateEnabled,maximumThroughputUnits:maximumThroughputUnits}' --output json > namespace-capacity.json

Des requêtes throttled ou des erreurs ServiceBusy associées à un débit au plafond justifient une action de capacité. L’absence de throttling avec un egress qui progresse faiblement oriente plutôt vers le consumer ou son aval. Une seule partition chaude oriente vers la clé de partition et la distribution ; ajouter des unités peut rendre de la marge, mais ne corrige pas une concentration durable.

La majoration automatique du niveau Standard augmente les unités jusqu’au maximum configuré, mais ne les réduit pas automatiquement. Notez donc l’état initial, le plafond et le coût attendu avant de le modifier. Ne présentez pas cette action comme un autoscaling complet.

Corréler retard, erreurs et déploiements

Les noms de tables dépendent du routage des diagnostics et de la télémétrie applicative. La requête suivante illustre la logique attendue : reconstruire par partition le dernier progrès, les erreurs et la version du consumer. Adaptez les tables et colonnes au schéma réellement envoyé dans Log Analytics.

kusto 02-consumer-partition-progress.kql
let Window = 2h;
let ConsumerGroup = "<consumer-group>";
ConsumerProgress_CL
| where TimeGenerated > ago(Window)
| where ConsumerGroup_s == ConsumerGroup
| summarize
  LastReceivedSequence = max(ReceivedSequence_d),
  LastProcessedSequence = max(ProcessedSequence_d),
  LastCheckpointSequence = max(CheckpointSequence_d),
  LastProgress = max(TimeGenerated),
  ProcessingP95Ms = percentile(ProcessingDurationMs_d, 95),
  Errors = countif(Result_s != "success")
by PartitionId_s, ConsumerVersion_s
| extend
  ProcessingGap = LastReceivedSequence - LastProcessedSequence,
  CheckpointGap = LastProcessedSequence - LastCheckpointSequence,
  ProgressAge = now() - LastProgress
| order by ProcessingGap desc

Ajoutez la latence et les erreurs de la dépendance aval dans la même fenêtre. Si ProcessingGap monte pendant que les appels aval ralentissent, scaler Event Hubs traite le mauvais composant. Si CheckpointGap monte alors que les succès métier continuent, suspendez tout replay : le risque principal est de retraiter des effets déjà appliqués.

Corrélez enfin la première divergence avec les changements : nouvelle version, taille de batch, fréquence de checkpoint, stratégie de retry, timeout aval, clé de partition, nombre d’instances, droits du checkpoint store ou règle réseau. Le dernier changement n’est pas automatiquement la cause, mais il fournit un rollback testable.

Séparer les familles de panne

text event-hubs-lag-classification.txt
Capacite du namespace atteinte
ThrottledRequests ou ServiceBusy augmente
Plusieurs flux du namespace sont touches
Action: capacite bornee ou auto-inflate, puis verifier le debit utile

Partition chaude
Une partition porte l'essentiel du retard
Les autres continuent de progresser
Action immediate: isoler le producteur ou reduire sa pression
Action durable: revoir la cle de partition et la distribution

Consumer sous-dimensionne
Partitions reparties mais duree de traitement en hausse
Instances saturees, sans erreur de service Event Hubs
Action: scaler jusqu'au nombre de partitions utiles et verifier l'aval

Dependance aval lente ou en erreur
Retries et latence aval montent avant le lag
Action: proteger l'aval, borner les retries, rollbacker le changement causal

Checkpoint bloque ou en retard
Traitement metier reussi, checkpoint stagnant
Action: reparer ownership, stockage ou permission sans remettre l'offset en arriere

Ownership instable
Partitions changent frequemment de consumer
Action: stabiliser instances, leases et timeouts avant tout scale supplementaire

Cette classification impose de traiter séparément capacité de transport et capacité de traitement. Elle évite aussi de confondre rétention et sauvegarde : Event Hubs permet de relire pendant la période conservée, mais ne garantit pas que l’application aval supporte les doublons.

Choisir une action bornée

Pour un throttling prouvé, augmentez temporairement la capacité ou le plafond d’auto-inflate avec une valeur maximale explicite, une alerte de coût et une heure de réévaluation. Pour un consumer lent, augmentez les instances seulement si des partitions restent sans owner ou si le parallélisme utile peut réellement progresser. Au-delà du nombre de partitions actives, des instances supplémentaires restent inactives ou provoquent davantage de rééquilibrages.

Pour une partition chaude, une correction de clé de partition ne redistribue pas l’historique existant. Stabilisez d’abord le producteur, laissez le consumer rattraper ou créez un flux de migration conçu et testé. Pour une dépendance aval, appliquez backpressure, circuit breaker ou réduction de batch selon le contrat ; évitez une vague de retries qui amplifie l’incident.

yaml bounded-recovery-decision.yml
decision:
diagnosis: <namespace-capacity|hot-partition|slow-consumer|downstream|checkpoint>
evidence_window: <start/end UTC>
target_consumer_group: <consumer-group>

action:
type: <scale-capacity|scale-consumer|rollback-release|repair-checkpoint-path|bounded-replay>
maximum_capacity: <value-or-not-applicable>
maximum_instances: <value-or-not-applicable>
affected_partitions: [<ids>]
owner: <operator>
expires_at_utc: <timestamp>

validation:
- no new throttled requests
- every partition has one stable owner
- processed sequence progresses on every partition
- checkpoint follows successful processing
- downstream error rate remains within baseline
- end-to-end delay decreases for three consecutive windows

rollback:
- restore previous consumer release and configuration
- restore previous capacity after backlog recovery
- stop replay if duplicate or downstream errors rise
- preserve checkpoint and processed-event evidence

Encadrer un replay sans dupliquer la production

Un replay est une nouvelle écriture métier, pas une simple opération de lecture. Avant de déplacer un offset, prouvez quelles partitions et quelles séquences manquent réellement. Vérifiez si le traitement est idempotent, si une clé de déduplication existe et si les effets partiels sont visibles.

Préférez un consumer group dédié au replay avec une plage temporelle ou des séquences bornées. Conservez le consumer group de production et ses checkpoints. Limitez le débit du replay pour protéger l’aval, journalisez chaque événement accepté, ignoré ou échoué et arrêtez sur un seuil d’erreur explicite.

text replay-gates.txt
Replay autorise seulement si
les partitions et bornes de sequence sont connues
la retention contient encore toute la plage
l'effet aval est idempotent ou dedupliquable
le consumer group de production reste intact
le debit maximal et le seuil d'arret sont definis
une preuve avant/apres existe pour chaque effet attendu

Arret immediat si
un evenement deja applique produit un doublon
le taux d'erreur aval depasse le seuil
le retard du flux de production recommence a monter
la borne ou l'ownership des checkpoints devient ambigu

Remettre directement le checkpoint de production en arrière mélange reprise et exploitation normale. Cette option ne doit être retenue que si le mécanisme de checkpoint est versionné, restaurable, revu et si les doublons sont maîtrisés. Dans la plupart des cas, un consumer group temporaire est plus observable et plus facile à arrêter.

Valider la reprise et rollbacker proprement

La reprise est prouvée lorsque chaque partition progresse, que l’âge du checkpoint diminue, que la latence aval reste saine et que le délai bout en bout revient vers le SLO sur plusieurs fenêtres. Un débit sortant élevé pendant quelques minutes ne suffit pas : il peut refléter un replay agressif ou un rééquilibrage.

Si la correction ne réduit pas le retard, revenez au dernier état stable : version précédente du consumer, paramètres de batch et de checkpoint précédents, nombre d’instances connu, puis capacité du namespace initiale après absorption du backlog. Ne réduisez pas les unités tant que le flux courant et le rattrapage partagent encore le même budget.

Documentez la cause au bon niveau. Une hausse de capacité résout une saturation agrégée ; elle ne corrige ni une partition chaude, ni une permission du checkpoint store, ni une base aval lente. La mesure de prévention doit suivre la classification : télémétrie par partition, alertes sur absence de progrès, test de replay, clé de partition revue ou budget d’auto-inflate.

Conclusion

Un retard Event Hubs se diagnostique comme une chaîne de progression : événement publié, partition attribuée, événement reçu, traitement validé, checkpoint écrit et effet aval confirmé. Tant qu’un maillon précis n’est pas identifié, scaler ou rejouer reste une hypothèse coûteuse.

La bonne décision est celle qui réduit l’écart sans perdre l’ownership ni dupliquer les effets : capacité bornée pour un throttling prouvé, parallélisme utile pour un consumer lent, correction de distribution pour une partition chaude, rollback applicatif pour une dépendance dégradée et replay isolé seulement quand les bornes et l’idempotence sont démontrées.