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
| Method | When called | What 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 completedonPartitionsRevoked()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
ConsumerRebalanceListenerto thesubscribe()method so it will get invoked."
4.7 Commits and offsets — where consumers actually break
Get this wrong by one and you reprocess one record per commit forever (harmless-ish) or skip one record per commit (data loss).
4.9 Consuming from specific offsets
① Map every partition assigned to this consumer (consumer.assignment()) to the target timestamp.