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.
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.
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.
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.
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.
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
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.
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.
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.