Messaging
Part 4 of 5 · Kafka messagingConsumer Rebalancing, Lag & Backpressure
Cooperative sticky vs eager; lag as offset gap; pause/resume; scale consumers vs partitions.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Why cooperative sticky plus pause beats deploy STW
Prefer
Cooperative sticky + static membership + pause
Revoke only partitions that change owners. group.instance.id survives a bounce. pause when the DB is saturated instead of OOM.
- Frequent deploys stop stalling the whole group.
- Lag on one partition is visible instead of averaged away.
- max.poll.interval.ms sized to real batch time.
- Idle extra pods do not magically drain a hot key.
Alternative
Eager rebalance, tiny poll interval, scale past N
Every deploy revokes everything. Slow batches get kicked. More consumers than partitions sit idle while one partition melts.
- Timeout storms: processing longer than max.poll.interval.ms.
- Averages hide a hot partition.
- Blocking in the poll loop without pause.
- Commit async then crash → surprise reprocessing (ALO).
Trigger, protocol, assign, then backpressure
Vertical cards for phones. Same path as the mermaid.
- 1
Trigger
Join, leave, subscribe change, or session/poll timeout. - 2
Protocol
Eager revokes ALL partitions. Cooperative revokes only partitions that must move. - 3
Assign and resume
New assignment. Poll, process, commit. Static membership skips this on a bounce. - 4
Lag rises
log end minus consumer offset per partition. Time-lag for freshness SLOs. - 5
Backpressure
consumer.pause, bound the worker queue, scale compute, or fix the hot key. Do not OOM.
Overview
A rebalance reassigns topic partitions among consumer-group members when membership or subscribed topics change. Eager protocols are stop-the-world (all pause, revoke all, reassign). Cooperative sticky revokes only what must move — less thrash. Lag = (high watermark − committed/consumer offset) per partition; it is a backlog signal, not a latency SLO by itself. Backpressure means slowing intake (pause/resume, bounded queues, fewer polls) instead of OOM. Scale consumers up to partition count; beyond that, add partitions (with affinity costs) or speed up handlers.
You should be able to:
- Sketch JoinGroup → revoke subset → SyncGroup.
- Size
session.timeout.ms, heartbeat,max.poll.interval.ms,max.poll.recordstogether. - Diagnose a deploy-induced rebalance loop.
Architecture
3 Assign
- 1
Revoke ALL partitions
- nextNew assignment
- 2
Revoke only moving partitions
- nextNew assignment
- 3
New assignment
- nextPoll / process / commit
4 Resume and backpressure
- 4
Poll / process / commit
- slowLag rises
- 5
Lag rises
- nextBackpressure
- 6
Backpressure
- pauseconsumer.pause partitions
- scaleMore consumers or faster code
- 7
consumer.pause partitions
- 8
More consumers or faster code
Flow
- 9
1 Trigger: join / leave / subscribe / timeout
- next2 Eager protocol?
- 10
2 Eager protocol?
- yesRevoke ALL partitions
- cooperativeRevoke only moving partitions
Lesson map
Consumer Rebalancing, Lag & Backpressure
Cooperative sticky vs eager; lag as offset gap; pause/resume; scale consumers vs partitions.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB m1["Consumer M1"] m2["Consumer M2"] coord["Group coordinator"] m2 -->|JoinGroup| coord coord -->|revoke| m1 coord -->|SyncGroup| m1 coord -->|SyncGroup| m2
Rebalance strategies
| Strategy | Stop-the-world? | Partition movement | Prefer when |
|---|---|---|---|
| Eager range/round-robin | Yes | Often many | Tiny groups, simple ops |
| Eager sticky | Yes | Fewer moves | Legacy clients |
| Cooperative sticky | Incremental | Minimal | Production default (modern clients) |
Static membership (group.instance.id) | Rare | Minimal on bounce | Stateful / large partitions |
What fails if you choose wrong
- Eager rebalance every deploy → processing stalls, timeout storms.
max.poll.interval.mstoo low for slow batches → consumer kicked mid-work → rebalance loops.- Ignoring lag on one hot partition while group lag looks fine (averages lie).
- Adding consumers without adding partitions → idle members, same lag. Hub rule: at most one member per partition.
Lag math, timeouts, backpressure
Lag math
lag(p) = log_end_offset(p) − consumer_offset(p). Sum of lags ≠ user latency. Track p99 partition lag and time-lag (event timestamp vs now) for freshness SLOs.
Commit-after-process is at-least-once: a kick mid-batch redelivers. Dedup/outbox still required: at-least-once vs exactly-once, outbox/inbox.
Critical timeouts
| Config | Role | Failure if wrong |
|---|---|---|
session.timeout.ms | Heartbeat liveness | Too low → flapping |
heartbeat.interval.ms | ≤ 1/3 session | Missed heartbeats |
max.poll.interval.ms | Max time between polls | Long processing → kick |
max.poll.records | Batch size | Huge batches stretch the interval |
Backpressure toolkit
consumer.pause(partitions)when downstream is saturated;resumewhen OK.- Bounded worker queue + stop polling when full.
- Slow consumers: split group, async handoff with care for ordering.
- Fix the hot key before “just scale”.
Sequence
- 1
Consumer M1
1 M2 joins
- 2
Consumer M2 → Group coordinator
JoinGroup
- 3
Group coordinator → Consumer M1
revoke cooperative subset
- 4
Consumer M1
2 Sync assignment
- 5
Group coordinator → Consumer M1
SyncGroup assignment
- 6
Group coordinator → Consumer M2
SyncGroup assignment
- 7
Consumer M1
3 Process and watch lag
- 8
Consumer M1 → Consumer M1
poll / process / commit
- 9
Consumer M2 → Consumer M2
poll / process / commit
Lag and pause (run this)
Downstream brownout pauses assigned partitions. Lag stays. Resume drains with bounded capacity.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Pros / cons
| Lever | Pros | Cons | Prefer when |
|---|---|---|---|
| More consumers | Parallelism | Idle if greater than partitions | Lag + CPU-bound handlers |
| Cooperative sticky | Less STW | Need modern clients | Frequent deploys |
| pause/resume | Protects DB | Lag grows intentionally | Downstream brownout |
| Static membership | Bounce-friendly | Session fencing care | K8s rolling deploys |
Interview Q&A
Max useful consumers in a group?
Answer
About the number of partitions subscribed (per topic assignment). Extra members idle.
Eager vs cooperative?
Answer
Eager revokes all. Cooperative revokes only partitions that change owners.
What is consumer lag?
Answer
Unread records per partition: log end − consumer offset. Not by itself a latency SLO.
Why rebalance loops?
Answer
Processing longer than max.poll.interval.ms, or session timeouts from GC/network. Deploys without static membership amplify it.
How do you apply backpressure?
Answer
pause partitions, bound in-flight work, scale compute, fix hot keys.
Is sum(lag) a good SLO?
Answer
Prefer age of oldest unprocessed event / time-lag for freshness.
Pitfalls
Six pods, 12 partitions, eager assignor, max.poll.interval.ms=30s, handler p99=40s. Rolling deploy one pod at a time. Write why members get kicked, why lag spikes, and what you change first: cooperative sticky, static membership, or the interval. Then add a 13th pod and mark it idle.