Kafka Consumer Groups Under Failure: What a Killed Consumer Costs with Range, Sticky, Cooperative and KIP-848 Assignment, Measured on Kafka 4.3.1

Lab run with apache/kafka:4.3.1 (KRaft), kafka-clients 4.3.1, Java 17 and confluent-kafka 2.15.1; 11 integration tests plus kafka-consumer-groups.sh output

Revision note (2026-09-15). The earlier version of this article called RangeAssignor "the default", which has not been the whole story since Kafka 3.0 and is not applicable at all under the KIP-848 protocol, where the client rejects partition.assignment.strategy and the broker picks the assignor. It dated incremental cooperative rebalancing for consumers to "Kafka 2.3+" (that release added it to Kafka Connect; the consumer got it in 2.4 with KIP-429, and the 2.3.1 client jar has no CooperativeStickyAssignor). Its code configured StickyAssignor, an eager assignor, while the text recommended cooperative rebalancing. It said to "commit offsets before partitions are revoked to avoid message duplication" without saying that the callback never runs when a consumer dies, and it never mentioned the session timeout that decides how long a dead consumer's partitions sit idle. Its printf strings had lost their newline escapes (offsetsn). This version is rebuilt around a lab with 11 integration tests on Kafka 4.3.1; every number below comes from that lab, and the lab's README has the raw output.

The lab

  • Broker: apache/kafka:4.3.1, single KRaft node in Docker. Client: kafka-clients 4.3.1 on Java 17 (Corretto 17.0.14). Python: confluent-kafka 2.15.1.
  • One topic, 6 partitions, 8,000 records per partition (48,000 total), prefilled before the consumers start.
  • Three consumers A, B, C, each in its own JVM so that "kill" means SIGKILL. Their loop is the earlier article's loop with logging added: subscribe with a ConsumerRebalanceListener whose onPartitionsRevoked calls commitSync(), poll, process, then commitAsync() after every batch. max.poll.records=50, 2 ms of simulated work per record.
  • Session timeout 6 s and heartbeat 2 s on both protocols (client settings for classic; group.consumer.session.timeout.ms and group.consumer.heartbeat.interval.ms on the broker for consumer). The defaults are 45 s and 3 s (classic) or 45 s and 5 s (KIP-848); the lab shortens them so a run takes seconds, not minutes.
  • Four assignment modes: RangeAssignor, StickyAssignor, CooperativeStickyAssignor (all group.protocol=classic), and group.protocol=consumer (KIP-848, server-side uniform assignor).
  • Two ways to remove C: SIGKILL, or a graceful consumer.close().

Every consumer writes each callback, each processed record and each commit acknowledgement with a timestamp. The test then counts records processed twice, checks that no committed offset hides an unprocessed record, and reads the group with the Admin API and kafka-consumer-groups.sh.

Where the assignor lives in Kafka 4.3

Two consumer group protocols exist side by side. group.protocol=classic is the default in kafka-clients 4.3.1; group.protocol=consumer is the KIP-848 protocol, GA since Kafka 4.0. The Kafka 4.3 documentation's evolution timeline (KIP-1274) says the client is expected to default to consumer in 5.0 and to drop classic in 6.0.

What the lab observed on 4.3.1:

partition.assignment.strategy default = [RangeAssignor, CooperativeStickyAssignor]
group.protocol default = classic
session.timeout.ms default = 45000, heartbeat.interval.ms default = 3000, max.poll.interval.ms default = 300000

The default strategy is a list, not RangeAssignor alone. RangeAssignor is what a fresh group uses; the second entry exists so that a rolling bounce which removes RangeAssignor from the list upgrades the group to cooperative rebalancing without downtime.

Under the KIP-848 protocol the whole client-side assignor concept goes away. With group.protocol=consumer, setting the strategy the earlier article recommended fails at construction time:

ConfigException: partition.assignment.strategy cannot be set when group.protocol=CONSUMER
ConfigException: session.timeout.ms cannot be set when group.protocol=CONSUMER

The broker's group.consumer.assignors (default uniform,range) decides, and the client can only pick one of the broker's assignors by name:

group.protocol=consumer without a client assignor -> type=CONSUMER state=Stable assignor=uniform
group.remote.assignor=range                       -> type=CONSUMER state=Stable assignor=range

The 4.3 documentation maps the old client-side assignors to the new server-side ones: RangeAssignor to range; CooperativeStickyAssignor, StickyAssignor and RoundRobinAssignor all to uniform. Custom client-side assignors (ConsumerPartitionAssignor) are not supported under the new protocol; the server-side extension point is ConsumerGroupPartitionAssignor in the broker configuration.

So the configuration block from the earlier article is protocol-dependent. For classic:

group.protocol=classic
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
session.timeout.ms=45000
heartbeat.interval.ms=3000

For KIP-848, none of those three lines is allowed:

group.protocol=consumer
# optional: group.remote.assignor=range   (default: the broker's first assignor, uniform)

What actually happens when a consumer dies

C is killed at t0 while all three are processing a large backlog. The survivors' logs give the moment they first heard about it, what was revoked from them, and how many records were processed twice.

ModeC ran onPartitionsRevokedSurvivors' first callbackRevoked from survivorsRecords processed twice
RangeAssignor, SIGKILLno+7.8 s / +7.9 severything they owned44
StickyAssignor, SIGKILLno+7.9 s / +7.8 severything they owned46
CooperativeStickyAssignor, SIGKILLno+8.0 s / +8.0 snothing19
KIP-848 (uniform), SIGKILLno+5.7 s / +7.6 snothing15
RangeAssignor, close()yes+1.8 s / +1.7 severything they owned0
CooperativeStickyAssignor, close()yes+1.9 s / +1.9 snothing0
KIP-848 (uniform), close()yes+1.7 s / +1.7 snothing0

In every run, zero offsets below the committed offset went unprocessed. The dangerous direction here is duplication, not loss.

Three things follow from this table.

A crash runs no callback. onPartitionsRevoked fired in all three close() runs and in none of the SIGKILL runs. The earlier article's advice to "commit offsets before partitions are revoked to avoid message duplication" describes the graceful path only. For a process that is killed, OOM-killed, or loses its network, there is no "before".

The wait is the session timeout, not the rebalance. With a 6 s session timeout and 2 s heartbeats, survivors first heard about C's death 5.7 to 8.0 s after the kill. With the 45 s defaults that wait is about 45 s (not measured here; the lab only ran the shortened values). During that time C's partitions are owned by nobody and lag grows on them. After a graceful close(), the LeaveGroup request made the survivors react within 1.7 to 1.9 s, inside one heartbeat interval (2 s in the lab). The rebalance itself, once it started, was fast: the longest gap a survivor saw between two processed records during an eager rebalance was 179 ms, against 14 ms in steady state.

Duplicates come from the commit lag, and the count depends on luck, not on the assignor. 44, 46, 19 and 15 are not properties of the four modes. Each equals exactly the number of records C had processed past the offset the broker had committed at the moment of the kill, and the test asserts that equality. The loop commits with commitAsync() after each batch of up to 50, so the window is "however far into the current batch C was, plus any commit still in flight". One raw line shows the in-flight part:

[lab range/kill9] C last acknowledged commit callback: {4=487, 5=337}; broker committed right after removal: {4=537, 5=337}; C last processed offset: {4=580, 5=336} => 44 records processed by C sit above the committed offset

C's log had acknowledged a commit for offset 487 on partition 4, but the broker already held 537: the request had arrived, the callback never ran. The survivor then re-read 537..580 (44 records). To shrink the window, commit more often or synchronously after each batch, and accept the throughput cost; to remove duplicates, make processing idempotent or use transactions. The lab measured neither of those variants.

Eager and cooperative differ in what survivors go through

The "Revoked from survivors" column is where the assignors differ. With RangeAssignor and StickyAssignor, both classic eager assignors, every survivor had its entire assignment revoked and then received a new one, and the article's commitSync() in onPartitionsRevoked ran on each of them for partitions that came straight back:

[lab sticky/kill9] before: A=[0,3], B=[1,4], C=[2,5]
[lab sticky/kill9] A: revoked [0, 3], revoke->assigned 12 ms, now owns [0, 2, 3]
[lab sticky/kill9] B: revoked [1, 4], revoke->assigned 77 ms, now owns [1, 4, 5]

StickyAssignor kept A on 0 and 3 and B on 1 and 4, which is what "sticky" promises. But it is still the eager protocol: the survivors gave everything up and got it back. Any per-partition state you tear down in onPartitionsRevoked is torn down for nothing.

With CooperativeStickyAssignor and with KIP-848, the survivors' onPartitionsRevoked was never called; onPartitionsAssigned delivered only the partition each survivor gained:

[lab cooperative/kill9] A: revoked [], now owns [0, 2, 3]
[lab cooperative/kill9] B: revoked [], now owns [1, 4, 5]

This is the practical meaning of "incremental": callbacks are about deltas. If your onPartitionsAssigned "reinitializes any local caches" for the whole assignment, as the earlier article suggested, it will throw away state for partitions that never moved. Initialize only what the callback hands you.

KIP-848 added one more visible difference. Its rebalance is per member, without a group-wide barrier: in the SIGKILL run, A received partition 5 at +5.7 s and B received partition 4 at +7.6 s, two separate reconciliations. The Python run showed the same shape on a two-consumer group: the first member got all six partitions, had three revoked one heartbeat later, and the second member received them the heartbeat after that.

Versions: 2.3, 2.4 and 4.0

  • Incremental cooperative rebalancing arrived in Kafka 2.3 for Kafka Connect workers (KIP-415, "Released: AK 2.3.0").
  • The consumer's CooperativeStickyAssignor arrived in Kafka 2.4 (KIP-429, "Accepted (2.4.0)"). kafka-clients-2.3.1.jar on Maven Central contains RangeAssignor, RoundRobinAssignor and StickyAssignor and no CooperativeStickyAssignor; kafka-clients-2.4.0.jar adds it together with the ConsumerPartitionAssignor interface.
  • The KIP-848 protocol is GA since Kafka 4.0 and must be enabled per client with group.protocol=consumer in 4.3.

The consumer loop, corrected

The structure of the earlier loop is fine; four details are not.

  1. The printf strings. The live post had "... committing offsetsn" and "... value = %sn". String.format with that text produces Partitions revoked: [] - committing offsetsn, no line break (pinned by a test). Use %n.
  2. commitSync() inside onPartitionsRevoked is correct for the graceful path and for eager rebalances, and harmless under cooperative rebalancing (the callback then runs only when a partition really leaves). Keep it, but do not describe it as protection against duplicates after a crash.
  3. Decide the protocol explicitly. Under classic, choose CooperativeStickyAssignor unless you have a reason not to. Under consumer, remove partition.assignment.strategy, session.timeout.ms and heartbeat.interval.ms from the client, and tune group.consumer.session.timeout.ms on the broker instead.
  4. Treat onPartitionsAssigned and onPartitionsRevoked as deltas, and implement onPartitionsLost for the case where the consumer discovers it has already been fenced (it is called instead of onPartitionsRevoked when ownership was lost rather than handed over; the default implementation delegates to onPartitionsRevoked).

The loop the lab ran, trimmed of its logging:

try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
    consumer.subscribe(List.of(topic), new ConsumerRebalanceListener() {
        @Override
        public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
            // Runs on close() and on eager rebalances; never after SIGKILL.
            consumer.commitSync();
        }

        @Override
        public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
            // Cooperative and KIP-848: only the partitions added in this rebalance.
        }

        @Override
        public void onPartitionsLost(Collection<TopicPartition> partitions) {
            // Ownership is already gone; committing here can fail.
        }
    });

    while (running) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
        for (ConsumerRecord<String, String> record : records) {
            process(record);
        }
        consumer.commitAsync((offsets, exception) -> {
            if (exception != null) {
                log.warn("commit failed: {}", exception.toString());
            }
        });
    }
}

With this loop, the duplicate window after a crash is at most one batch plus one in-flight commit, as measured above. Whether that is acceptable is a property of your processing, not of Kafka.

The Python snippet

The earlier article's confluent-kafka example runs as written on confluent-kafka 2.15.1 against Kafka 4.3.1. The lab ran it with two consumers and a graceful close of one of them, under range, cooperative-sticky, and group.protocol=consumer:

$ python rebalance_demo.py cooperative
+ 0.28s B on_assign  [1, 3, 5]
+ 0.28s A on_assign  [0, 2, 4]
+ 6.50s B on_revoke  [1, 3, 5]
+ 9.36s A on_assign  [5, 3, 1]
+12.69s committed offsets [(0, 100), (1, 100), (2, 100), (3, 100), (4, 100), (5, 100)]

Two notes for that client. Consumer.commit() with no arguments is asynchronous by default (asynchronous=True in the 2.15.1 docstring), so the commit in on_revoke is a request, not a completed write; pass asynchronous=False if you need to know it landed before the partitions move. And the same crash caveat applies: on_revoke runs on close() and on rebalances, not when the process is killed. The lab did not kill a Python consumer, so the duplicate window on that client was not measured.

What this does not cover

  • The 45 s default session timeout. The lab used 6 s on both protocols; the "about 45 s" above is the default value from the configuration reference, not a measurement.
  • max.poll.interval.ms eviction of a slow consumer, static membership (group.instance.id), custom assignors, multi-broker clusters, transactions and exactly-once.
  • Throughput. The 2 ms per record and a 4 KB per-partition fetch exist to make rebalances observable.
  • Monitoring stacks. Lag was read from kafka-consumer-groups.sh --describe and the Admin API only; no exporter, Cruise Control or dashboard was set up.

Reproduce it

cd examples/kafka-consumer-groups
./start-broker.sh                        # apache/kafka:4.3.1 on localhost:9096
gradle test --no-daemon --rerun-tasks    # 11 tests, about 2 min 15 s
./stop-broker.sh

The per-consumer logs of each scenario stay in build/scenario-logs/.

Sources