Learn Labs
4. Kafka Consumers: Reading Data from Kafka

4.8 Rebalance listeners

Pass a ConsumerRebalanceListener to subscribe().

Why: "a consumer will want to do some cleanup work before exiting and also before partition rebalancing. If you know your consumer is about to lose ownership of a partition, you will want to commit offsets of the last event you've processed. Perhaps you also need to close file handles, database connections, and such."

Pass a ConsumerRebalanceListener to subscribe().

The three methods

MethodWhen calledWhat to do there
onPartitionsAssigned(partitions)After partitions are reassigned to the consumer, before it starts consuming"prepare or load any state that you want to use with the partition, seek to the correct offsets if needed." ⚠️ "Any preparation done here should be guaranteed to return within max.poll.timeout.ms so the consumer can successfully join the group."
onPartitionsRevoked(partitions)When the consumer must give up partitions it owned — rebalance or close. Eager: invoked before rebalancing starts and after the consumer stopped consuming. Cooperative: invoked at the END of the rebalance, with just the subset being given up"This is where you want to commit offsets, so whoever gets this partition next will know where to start."
onPartitionsLost(partitions)Cooperative only, and only in exceptional cases where partitions were assigned to other consumers without first being revoked"clean up any state or resources used with these partitions." ⚠️ "has to be done carefully — the new owner may have already saved its own state, and you'll need to avoid conflicts." If you don't implement it, onPartitionsRevoked() is called instead

TIP — cooperative rebalancing callback semantics

  • onPartitionsAssigned() — invoked on every rebalance, as a way of notifying the consumer that a rebalance happened. If there are no new partitions assigned, it is called with an empty collection.
  • onPartitionsRevoked() — invoked in normal rebalancing conditions, but only if the consumer gave up ownership of partitions. It will NOT be called with an empty collection.
  • onPartitionsLost() — invoked in exceptional rebalancing conditions; the partitions in the collection will already have new owners by the time the method is invoked.

The ordering guarantee: "If you implemented all three methods, you are guaranteed that during a normal rebalance, onPartitionsAssigned() will be called by the new owner of the partitions only AFTER the previous owner completed onPartitionsRevoked() and gave up its ownership."

That guarantee is what makes handoff of external state (locks, file handles, DB transactions) actually safe.

The canonical example

private Map<TopicPartition, OffsetAndMetadata> currentOffsets = new HashMap<>();
Duration timeout = Duration.ofMillis(100);

private class HandleRebalance implements ConsumerRebalanceListener {
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
        // nothing needed — we just start consuming
    }

    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
        System.out.println("Lost partitions in rebalance. " +
            "Committing current offsets:" + currentOffsets);
        consumer.commitSync(currentOffsets);     // ← SYNC, on purpose
    }
}

try {
    consumer.subscribe(topics, new HandleRebalance());   // ← THE KEY LINE

    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(timeout);
        for (ConsumerRecord<String, String> record : records) {
            System.out.printf("topic = %s, partition = %s, offset = %d, " +
                 "customer = %s, country = %s\n",
                 record.topic(), record.partition(), record.offset(),
                 record.key(), record.value());
            currentOffsets.put(
                new TopicPartition(record.topic(), record.partition()),
                new OffsetAndMetadata(record.offset()+1, null));
        }
        consumer.commitAsync(currentOffsets, null);
    }
} catch (WakeupException e) {
    // ignore, we're closing
} catch (Exception e) {
    log.error("Unexpected error", e);
} finally {
    try {
        consumer.commitSync(currentOffsets);
    } finally {
        consumer.close();
        System.out.println("Closed consumer and we are done");
    }
}

Two authorial notes worth keeping:

  • "We are committing offsets for all partitions, not just the partitions we are about to lose — because the offsets are for events that were already processed, there is no harm in that."
  • "we are using commitSync() to make sure the offsets are committed before the rebalance proceeds."
  • "The most important part: pass the ConsumerRebalanceListener to the subscribe() method so it will get invoked."

On this page