skip to content

How do you diagnose partition skew from lag data, and how would you design lag-driven autoscaling for a consumer group?

level: principalimportance: should knowfreq 35%

answer

  1. Skew = lag on a few partitions, others ~0
  2. Hot key vs stuck/slow consumer
  3. More consumers can't fix skew — repartition
  4. Max replicas = partition count
  5. Scale on time-lag with hysteresis; rebalance spikes cause flapping

basics

~20 s

Skew shows as lag concentrated on a few partitions while others sit at zero — usually hot keys or a stuck/slow consumer. For autoscaling, scale consumers on total/time lag, but cap replicas at the partition count, throttle on rebalances, and prefer time-lag with hysteresis to avoid flapping.

solid answer

~50 s

Partition skew appears in per-partition lag as a few partitions with large, growing lag while the rest are near zero. Causes: uneven key distribution (a hot key pins traffic to one partition), a slow/stuck consumer owning those partitions, or a downstream dependency slow only for certain keys. You confirm with kafka-consumer-groups.sh --describe (per-partition LAG) and by correlating with throughput. For lag-driven autoscaling, expose lag to an autoscaler (KEDA's Kafka scaler, or Prometheus + HPA via Kafka Lag Exporter's lag_seconds). Critical design rules: (1) max consumers = partition count — extra replicas idle, so scaling is gated by partitioning; (2) scale on time-lag, not raw records, so the signal is throughput-normalized; (3) add hysteresis / cooldowns because every scale event triggers a rebalance that pauses consumption and briefly spikes lag; (4) skew can't be fixed by adding consumers — repartition or fix key distribution instead. Combine in-app records-lag-max for fast signal with external time-lag for correctness.

go deeper

for a junior

Recognize that skew means some partitions have much more lag than others.

for a middle

Identify hot-key vs slow-consumer causes and know lag can drive scaling decisions.

for a senior

Explain the partition-count ceiling, rebalance-induced spikes, and why time-lag beats record-lag for scaling.

for a principal

Architect the autoscaling system end to end — signal choice, replica ceiling, hysteresis, skew detection routed to repartitioning, dead-consumer safety — and set partition counts for peak fan-out.

## Diagnosing partition skew from lag **Skew** = work is unevenly distributed across partitions, so lag piles up on a subset. Symptom in `kafka-consumer-groups.sh --describe`: a handful of partitions show large and growing `LAG` while the rest are ~0. Distinguish two root causes: 1. **Data/key skew**: producers route by key hash; a **hot key** (or few keys) lands disproportionate traffic on one partition. That partition's LEO grows fast, so even a healthy consumer lags there. Look at per-partition *production rate*, not just lag. 2. **Consumer skew**: the partition's *production* is normal but its *owner* consumer is slow or stuck (GC, a poison record, a slow downstream call for those keys, an unbalanced assignment). Here LEO is normal but the committed offset isn't advancing — correlate lag with each consumer's processing rate and the `CONSUMER-ID` column. A **stuck single partition** (lag climbing on exactly one partition, owner unchanged, others fine) is a classic poison-message or downstream-deadlock signature. **Key insight:** adding more consumers does **not** fix skew. Kafka assigns whole partitions to consumers; one partition is processed by exactly one consumer in the group. If one partition is hot, more consumers just leaves them idle while the hot partition stays bottlenecked. Fixes are **repartitioning** (more partitions + better key), **changing the partition key** to spread hot keys, or **decoupling** heavy work (e.g. a second topic / parallelism within the consumer). ## Designing lag-driven autoscaling Goal: scale the number of consumer instances up when the group falls behind and down when caught up. **Signal source** - Prometheus + **Kafka Lag Exporter**'s `kafka_consumergroup_group_lag_seconds` → HPA, or - **KEDA** with its built-in **Kafka scaler** (triggers on lag threshold per partition / total), or - in-app `records-lag-max` for the fastest reaction (but remember it goes silent on crash). **Hard constraints and pitfalls** 1. **Replica ceiling = partition count.** A consumer group can have at most one active consumer per partition; replicas beyond the partition count sit idle. So autoscaling is bounded — you must partition for your peak parallelism up front (partition count is hard to decrease). 2. **Every scale event causes a rebalance.** Adding/removing a member triggers partition reassignment, during which consumption pauses (stop-the-world rebalance) and lag *spikes momentarily*. Naive autoscaling sees the spike and scales again → **flapping**. Mitigate with **cooldown/stabilization windows**, **scale-down delays**, and **hysteresis** (different up/down thresholds). Cooperative/incremental rebalancing (`CooperativeStickyAssignor`) reduces the pause but doesn't eliminate it. 3. **Use time-lag, not record-lag.** A fixed record threshold misfires across throughput regimes; `*_lag_seconds` maps directly to an end-to-end latency SLO and scales sensibly. 4. **Skew defeats horizontal scaling.** If lag is concentrated on a few partitions, more replicas won't help — detect skew first and route the response to repartitioning, not autoscaling. 5. **Scale-down safety.** Removing consumers also rebalances; drain in-flight work and avoid removing instances while lag is non-trivial. 6. **Avoid the dead-consumer trap.** If you trigger only on in-app metrics and a consumer dies, the metric vanishes; pair with external lag so 'no signal + growing broker lag' still triggers action/alerts. ## Putting it together A robust setup: KEDA/HPA driven by **time-lag** from an external exporter; **min replicas ≥ 1**, **max replicas = partition count**; **cooldown** windows to absorb rebalance spikes; a **separate alert** on per-partition skew (lag concentrated) that pages humans rather than autoscaling; and `CooperativeStickyAssignor` to soften rebalances. Partition count chosen at design time for peak fan-out.

  • You see lag growing on exactly one partition while the rest stay at zero. What are the likely causes and the fix?
    Either a hot key sending disproportionate traffic to that partition, or a slow/stuck consumer (poison message, slow downstream) owning it. Adding consumers won't help; repartition / fix the key, or fix the stuck processing.
  • Why can't you scale a consumer group's parallelism beyond the partition count?
    Within a group each partition is owned by exactly one consumer. Once there is one consumer per partition, additional instances get no assignment and idle. Partition count is the parallelism ceiling.
  • How do you stop lag-driven autoscaling from flapping?
    Each scale event triggers a rebalance that briefly spikes lag. Use cooldown/stabilization windows, hysteresis (separate up/down thresholds), scale-down delays, and cooperative rebalancing to dampen the oscillation.

saying these in an interview costs you the question

  • Claiming you can fix partition skew by adding more consumer instances.
  • Scaling beyond the partition count expecting more throughput.
  • Autoscaling on raw record-lag without hysteresis (flaps on rebalance spikes).
  • Ignoring that each scale event triggers a rebalance that itself spikes lag.
  • Driving autoscaling only on in-app metrics that vanish when consumers crash.

context