Cloud

Azure Event Hubs: diagnose consumer lag before scaling or replaying

A production runbook to separate ingestion pressure, throttling, hot partitions, consumer failures and checkpoint drift before adding capacity or replaying events.

13 Aug 2026 azureevent-hubsstreamingconsumer-lagcheckpointpartitionsazure-monitorkqlobservabilityautomationrunbookrollbackproduction

An Azure Event Hubs consumer is falling behind. Events still arrive, but business processing shows up several minutes late and the gap keeps growing. Adding Throughput Units, starting more instances or moving the checkpoint backward can all look like safe responses. Each can make the incident worse: extra consumers cannot create useful parallelism beyond the active partitions, more Event Hubs capacity cannot repair a slow downstream dependency, and an unbounded replay can duplicate effects that were already committed.

The running example is a near-real-time detection pipeline. A production consumer group processes events by partition, checkpoints after successful work and writes to a downstream database. This runbook drives an evidence-based choice: repair the consumer, remove namespace throttling, correct load distribution, run a bounded replay or roll back the latest change.

Freeze the symptom and scope

Name the exact flow first. Event Hubs capacity is shared at namespace level, while ownership and checkpoints belong to an event hub, consumer group and partition. An aggregate chart can therefore hide a neighboring producer consuming capacity or a single partition carrying most of the lag.

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

Freeze deployments, application autoscaling and manual checkpoint operations while collecting the first evidence set. Record the remaining retention window as well. When the oldest unprocessed event is approaching expiry, the decision window becomes part of the incident.

Build a partition evidence matrix

Lag is not one number. For every partition, compare the latest published sequence, the last successfully processed sequence, the durable checkpoint, the current owner and the age of the last progress. A consumer that reads quickly but checkpoints infrequently has a different failure from one that no longer receives events.

text partition-evidence-matrix.txt
For each partition
Event Hubs
  latest available sequence or offset
  observed ingress and egress rate
  errors or throttling in the incident window

Consumer
  current owner and latest lease acquisition
  latest received sequence
  latest successfully processed sequence
  p50 p95 p99 duration and processing errors

Checkpoint
  durable sequence or offset
  last update time
  gap from the latest successful processing

Downstream
  latency, errors, saturation and limits
  idempotency key or evidence of an applied effect

Classification
  uniform lag across partitions
  one hot partition
  active reads but slow processing
  successful processing but stalled checkpoint
  unstable ownership or unassigned partitions

Premium and Dedicated application metrics logs can expose signals such as ConsumerLag and OffsetCommit. In other configurations, or when that telemetry is not enabled, instrument the consumer itself: received sequence, committed business result, checkpoint write and processing duration. Do not derive precise lag from IncomingMessages - OutgoingMessages alone. Multiple consumer groups can read the same stream, and those counters are aggregated.

Prove or rule out a capacity limit

Inspect the namespace before blaming the application: IncomingMessages, OutgoingMessages, IncomingBytes, OutgoingBytes, ThrottledRequests, ServerErrors and UserErrors. Compare the incident window with a healthy baseline at the finest useful granularity. Standard-tier Throughput Units are shared by every event hub in the namespace, so another workload can be causal.

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

Throttled requests or ServiceBusy errors aligned with a capacity ceiling justify a capacity action. No throttling and weak egress progress point toward the consumer or its dependency. A single hot partition points toward the partition key and distribution strategy. Extra units may create headroom, but they do not repair persistent skew.

Standard-tier Auto-inflate raises units up to the configured maximum, but it does not automatically scale them down. Record the initial setting, ceiling and expected cost before changing it. Treat it as bounded scale-up, not full autoscaling.

Correlate lag, errors and deployments

Table names depend on how diagnostics and application telemetry are routed. The following query shows the required logic: reconstruct the last progress, errors and consumer version for each partition. Adapt the tables and columns to the schema actually sent to 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

Add downstream latency and failures for the same window. If ProcessingGap rises while downstream calls slow down, scaling Event Hubs targets the wrong component. If CheckpointGap rises while business operations keep succeeding, stop any replay plan: the primary risk is processing effects twice.

Correlate the first divergence with changes to release version, batch size, checkpoint frequency, retry strategy, downstream timeout, partition key, instance count, checkpoint-store permissions or network policy. The latest change is not automatically the cause, but it gives the team a testable rollback candidate.

Separate the failure families

text event-hubs-lag-classification.txt
Namespace capacity reached
ThrottledRequests or ServiceBusy increases
Multiple namespace workloads are affected
Action: bounded capacity or Auto-inflate, then verify useful throughput

Hot partition
One partition carries most of the lag
Other partitions continue to progress
Immediate action: isolate the producer or reduce its pressure
Durable action: review partition key and distribution

Undersized consumer
Partitions are balanced but processing duration rises
Instances are saturated without Event Hubs service errors
Action: scale up to useful partition parallelism and verify downstream health

Slow or failing downstream dependency
Downstream retries and latency rise before lag
Action: protect the dependency, bound retries, roll back the causal change

Stalled or delayed checkpoint
Business processing succeeds while checkpoint stops
Action: repair ownership, storage or permission without rewinding offsets

Unstable ownership
Partitions move repeatedly between consumers
Action: stabilize instances, leases and timeouts before adding more scale

This classification keeps transport capacity separate from processing capacity. It also prevents the team from treating retention as a backup. Event Hubs supports reading retained events again, but that does not guarantee the downstream system can accept duplicate effects.

Choose a bounded action

For proven throttling, temporarily raise capacity or the Auto-inflate ceiling with an explicit maximum, cost alert and review time. For a slow consumer, add instances only when partitions are unowned or useful parallelism can increase. Beyond the active partition count, more instances sit idle or trigger additional rebalancing.

For a hot partition, changing the partition key does not redistribute existing history. Stabilize the producer first, let the consumer catch up or design and test a migration stream. For a slow downstream service, apply backpressure, a circuit breaker or smaller batches according to its contract. Avoid a retry wave that amplifies the 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 the 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

Bound replay without duplicating production

A replay is a new business write, not merely another read. Before moving any offset, prove which partitions and sequences are actually missing. Verify idempotency, the deduplication key and visibility of partial effects.

Prefer a dedicated replay consumer group with a bounded time or sequence range. Preserve the production consumer group and its checkpoints. Rate-limit replay to protect downstream systems, log every accepted, skipped and failed event, and stop at an explicit error threshold.

text replay-gates.txt
Replay is allowed only when
partition and sequence bounds are known
retention still covers the complete range
downstream effects are idempotent or deduplicated
the production consumer group remains unchanged
maximum rate and stop threshold are defined
before and after evidence exists for every expected effect

Stop immediately when
an already applied event creates a duplicate
downstream error rate exceeds the threshold
production-stream lag starts rising again
checkpoint bounds or ownership become ambiguous

Rewinding the production checkpoint directly mixes recovery with normal operations. Use that path only when checkpoint state is versioned, restorable and reviewed, and duplicate effects are controlled. A temporary consumer group is usually easier to observe and stop.

Validate recovery and roll back cleanly

Recovery is proven when every partition progresses, checkpoint age falls, downstream latency stays healthy and end-to-end delay moves toward the SLO for several consecutive windows. A short burst of high egress is not sufficient evidence. It may only represent an aggressive replay or ownership rebalance.

If the change does not reduce lag, return to the last stable state: previous consumer release, batch and checkpoint settings, known instance count, then the initial namespace capacity after the backlog is absorbed. Do not scale units down while current traffic and recovery still share the same budget.

Document the cause at the right layer. Capacity fixes aggregate saturation; it does not repair a hot partition, checkpoint-store permission or slow database. Prevention should match the classification: per-partition progress telemetry, no-progress alerts, a tested replay procedure, reviewed partition keys or an Auto-inflate budget.

Conclusion

Event Hubs lag is best diagnosed as a progress chain: event published, partition owned, event received, processing committed, checkpoint written and downstream effect confirmed. Until one link is identified, scaling or replaying remains an expensive hypothesis.

The right decision reduces the gap without losing ownership or duplicating effects: bounded capacity for proven throttling, useful parallelism for a slow consumer, distribution changes for a hot partition, application rollback for a degraded dependency and an isolated replay only when sequence bounds and idempotency are demonstrated.