Rebalancing Kafka Consumer Groups Without Stalls

A Kafka consumer group redistributes partitions every time a member joins, leaves, or is suspected dead, and with default settings each redistribution pauses the entire group. This guide removes most of those pauses and shortens the rest, as part of Queue Partitioning Strategies in Queue Fundamentals & Architecture.

Problem Statement

A job-processing service consumes a 48-partition topic with 24 pods. Every deploy triggers a cascade of rebalances as pods restart one at a time; during each, all consumers stop processing, and the deploy takes lag from near zero to 90 seconds. Worse, a handful of slow jobs (occasionally 6 minutes) exceed max.poll.interval.ms, the broker evicts the consumer, the group rebalances, the slow job's partition moves to another pod — which reprocesses the batch — and the evicted pod rejoins, causing another rebalance. On bad days the group spends more time rebalancing than processing. You want deploys that move only the partitions that must move, no stop-the-world pauses, no evictions from slow jobs, and no duplicate processing on revocation.

Prerequisites

  • Kafka 2.4+ brokers and clients (cooperative rebalancing), 3.x recommended.
  • Consumers in a language with a mature client (Java, librdkafka-based Python/Go/.NET).
  • Deployments with stable pod identities if you want static membership (Kubernetes StatefulSet, or pod names passed as config).
  • Metrics for rebalance counts and consumer lag per partition.

Step 1 — Understand Why the Default Rebalance Stops Everything

With the classic eager protocol, a rebalance begins with every member revoking all its partitions, then the group leader computes a new assignment, then everyone resumes. Even if only one member left, 24 pods stop work, commit, and wait. The more members and partitions, the longer the pause.

eager rebalance (range / round-robin assignors):
  t0  member 7 leaves
  t1  ALL 24 members revoke ALL partitions -> processing stops everywhere
  t2  JoinGroup / SyncGroup round trips
  t3  new assignment; everyone resumes (most get the same partitions back)

Cooperative rebalancing changes step t1: members keep partitions that are not moving, and only the partitions that actually change owner are revoked, in a second, smaller round.

Stop-the-world vs incremental In an eager rebalance, when one member leaves, all members revoke all partitions and processing stops across the whole group until the new assignment is distributed. In a cooperative rebalance, members keep the partitions that are not moving and continue processing them; only the departed member's partitions are reassigned, so only those pause briefly. One member leaves a 24-member group eager all 48 partitions paused until the new assignment lands cooperative 46 partitions keep processing 2 Only the departed member's two partitions pause while they move to new owners.

Step 2 — Switch to the Cooperative-Sticky Assignor

The cooperative-sticky assignor keeps assignments as stable as possible and uses the incremental protocol. Switching a running group from an eager assignor needs a two-step rollout, because all members must support the protocol the group uses.

# Step A: deploy with both, eager first — the group keeps using the eager protocol
partition.assignment.strategy=range,cooperative-sticky

# Step B: after every member runs the Step A config, deploy with cooperative only
partition.assignment.strategy=cooperative-sticky
# confluent-kafka (librdkafka) equivalent
consumer = Consumer({
    "bootstrap.servers": BOOTSTRAP,
    "group.id": "jobs",
    "partition.assignment.strategy": "cooperative-sticky",
    "enable.auto.commit": False,
})

Skipping step A — deploying cooperative-sticky alone to a group running range — leaves members that cannot agree on a protocol, and new members fail to join until the old ones are gone.

Step 3 — Use Static Membership for Restarts

Even with cooperative rebalancing, a pod restart is a leave plus a join: two rebalances, and the partitions move away and back. Static membership gives each consumer a stable group.instance.id; when a member with a known id disconnects, the coordinator waits up to session.timeout.ms for it to return before rebalancing, and a quick restart rejoins with the same assignment and no rebalance at all.

import os
consumer = Consumer({
    "bootstrap.servers": BOOTSTRAP,
    "group.id": "jobs",
    "group.instance.id": os.environ["POD_NAME"],     # jobs-worker-0 ... stable across restarts
    "session.timeout.ms": 45000,                     # restart must complete within this
    "heartbeat.interval.ms": 3000,
    "partition.assignment.strategy": "cooperative-sticky",
})
Restarts without rebalances With dynamic membership, a restarting pod leaves the group, triggering a rebalance that moves its partitions away, and rejoins, triggering a second rebalance that may move them back. With static membership, the coordinator holds the pod's partitions for up to the session timeout; the restarted pod rejoins with the same instance id and resumes its partitions with no rebalance. Restarting jobs-worker-7 dynamic leave: rebalance 1 join: rebalance 2 static partitions held for worker-7 (under 45 s) rejoin: same assignment A rolling deploy of 24 pods goes from 48 rebalances to zero.

A StatefulSet gives stable pod names; with a Deployment, derive a stable id from something else (a slot number) — random pod names defeat static membership. Size session.timeout.ms above the time a pod takes to restart; the trade-off is that a genuinely dead static member keeps its partitions unassigned for that long.

Step 4 — Stop Slow Jobs from Evicting Consumers

max.poll.interval.ms is the maximum time between poll() calls. A consumer that processes a batch longer than that is considered stuck and is removed from the group — the eviction cascade in the problem statement. Two fixes, in order of preference:

# 1. Decouple processing from polling: keep polling (paused) while slow work runs
consumer.pause(consumer.assignment())                # stop fetching new records
future = executor.submit(process_batch, records)
while not future.done():
    consumer.poll(0.5)                               # heartbeats + keeps the member alive
consumer.resume(consumer.assignment())

# 2. And bound the batch so a poll loop turn is short
consumer_config.update({
    "max.poll.interval.ms": 600000,                  # 10 min ceiling as a backstop
})
records = consumer.consume(num_messages=50, timeout=1.0)   # small batches for slow jobs

Pausing and polling keeps the member alive during long work without fetching more data. Raising max.poll.interval.ms alone just delays detection of genuinely stuck consumers. Jobs that routinely take minutes are often better served by a job queue with per-message acknowledgement — see Kafka vs RabbitMQ for task queues.

Keep polling while the slow job runs Without pausing, the consumer does not call poll for six minutes, exceeding the five-minute max poll interval; the coordinator evicts it, a rebalance moves its partitions, and the batch is reprocessed elsewhere. With pause and poll, the consumer pauses its partitions, runs the job in a thread, and keeps polling every half second, so it stays in the group and no rebalance happens. A 6-minute job, max.poll.interval.ms = 5 minutes no poll processing, no poll() evicted, rebalance pause + poll job in a thread, poll(0.5) keeps the member alive 5-minute limit Paused partitions fetch nothing, so memory stays bounded while the job runs.

Step 5 — Commit and Clean Up on Revocation

When partitions are revoked, commit what has been processed for them before they move, so the new owner starts from the right offset. With cooperative rebalancing, the revoke callback receives only the partitions being taken away.

def on_revoke(consumer, partitions):
    # finish or abandon in-flight work for these partitions only
    for tp in partitions:
        wait_for_in_flight(tp, timeout=10)
    offsets = [TopicPartition(tp.topic, tp.partition, next_offset(tp)) for tp in partitions]
    if offsets:
        consumer.commit(offsets=offsets, asynchronous=False)
    for tp in partitions:
        drop_local_state(tp)

def on_assign(consumer, partitions):
    for tp in partitions:
        init_local_state(tp)                     # e.g., per-partition caches, dedup windows

consumer.subscribe(["jobs"], on_assign=on_assign, on_revoke=on_revoke)

The synchronous commit in on_revoke is what prevents the reprocessing burst after each rebalance; without it, the new owner restarts from the last periodic commit. Handlers must still be idempotent, because a member that dies without running on_revoke leaves uncommitted work. The ordering implications are covered in per-entity ordering with Kafka partition keys.

Step 6 — Make Deploys Rebalance-Aware

With static membership and cooperative assignment, a rolling restart that keeps each pod's downtime under session.timeout.ms causes no rebalances. Configure the rollout to match:

spec:
  updateStrategy:
    type: RollingUpdate
    rollingUpdate: { partition: 0 }       # StatefulSet: one pod at a time, in order
  template:
    spec:
      terminationGracePeriodSeconds: 30   # close() commits and leaves cleanly within this
      containers:
        - name: worker
          readinessProbe:                 # "ready" only after partitions are assigned
            exec: { command: ["/bin/sh", "-c", "test -f /tmp/assigned"] }
            periodSeconds: 2

Gating readiness on partition assignment makes the rollout wait for each pod to rejoin before restarting the next, which keeps at most one pod's partitions in motion at a time. Note that a static member calling close() does not leave the group (so its partitions wait for it), which is exactly what you want for restarts and not what you want when scaling down — scale down by removing pods whose ids will not return, and accept one rebalance after the session timeout.

Verification

# Rebalances per hour should drop to near zero outside scaling events
kafka-consumer-groups.sh --bootstrap-server $B --describe --group jobs --state
# Client metrics (Java): rebalance-rate-per-hour, rebalance-latency-avg, failed-rebalance-total

Run a rolling restart of all 24 pods under normal load and compare maximum partition lag before and after the changes; in the scenario it fell from 90 seconds of lag per deploy to under 3 seconds. Then run a six-minute job and confirm no eviction appears in the coordinator logs.

Gotchas & Edge Cases

Mixed client versions. Old clients that do not support cooperative rebalancing force the group to eager behaviour. Upgrade all consumers before switching.

Static ids reused concurrently. Two running pods with the same group.instance.id fence each other repeatedly. Ensure ids are unique per running instance.

Too many partitions per member. Each partition adds work to rebalances and commits. Partition count is a long-term decision; see Queue Partitioning Strategies.

Autoscaling churn. An autoscaler that adds and removes consumers every few minutes causes a rebalance each time. Use scale-down stabilisation windows of several minutes.

FAQ

Does KIP-848 (the new consumer rebalance protocol) change this? Kafka 4.x's server-side protocol (group.protocol=consumer) moves assignment to the broker and makes rebalances incremental by design, removing much of the tuning here. Where available, adopt it; the timeout and slow-job advice still applies.

How do I see how much time rebalances cost? The Java client exposes rebalance-latency-avg, rebalance-total, and last-rebalance-seconds-ago per consumer; librdkafka reports rebalance events through statistics callbacks. Combine them with per-partition lag: a sawtooth in lag that lines up with deploys or autoscaling events is the rebalance cost made visible, and it should flatten after the changes above.

What session timeout should I use? Long enough to cover a pod restart with static membership (30–60 s), short enough that a dead member's partitions do not sit idle too long. Heartbeat interval about one-third of it or less.

Why do I still see rebalances with cooperative-sticky? Members still join and leave on scaling, crashes, and evictions. Cooperative makes each rebalance cheap; static membership and slow-job handling reduce how many happen.

Related