5.6 Consumer group management
valid() is the right choice for tooling that must degrade gracefully; all() is the right choice when partial results would be misleading.
Why programmatic offset control matters:
"unlike most message queues, Kafka allows you to reprocess data in the exact order in which it was consumed and processed earlier. In Chapter 4 we explained how to use the Consumer APIs to go back and reread older messages. But using these APIs means that you programmed the ability to reprocess data in advance into your application. Your application itself must expose the 'reprocess' functionality."
Two scenarios where you need it and the app doesn't provide it:
- Troubleshooting a malfunctioning application during an incident
- Preparing an application to start running on a new cluster during a disaster recovery failover
6.1 List consumer groups — and the valid() / errors() / all() distinction
admin.listConsumerGroups().valid().get().forEach(System.out::println);| Call | What you get |
|---|---|
.valid() | Only groups the cluster returned without errors. “Any errors will be completely ignored, rather than thrown as exceptions.” |
.errors() | Get all the exceptions. |
.all() | “Only the first error the cluster returned will be thrown as an exception.” |
Likely causes of such errors: "authorization, where you don't have permission to view the group, or cases when the coordinator for some of the consumer groups is not available."
valid() is the right choice for tooling that must degrade gracefully; all() is the right choice when partial results would be misleading.
6.2 Describe a group
ConsumerGroupDescription groupDescription = admin
.describeConsumerGroups(CONSUMER_GRP_LIST)
.describedGroups().get(CONSUMER_GROUP).get();
System.out.println("Description of group " + CONSUMER_GROUP
+ ":" + groupDescription);What you get:
What you get:
- the group members
- their identifiers and hosts
- the partitions assigned to them
- the algorithm used for the assignment
- the host of the group coordinator
"This description is very useful when troubleshooting consumer groups. One of the most important pieces of information about a consumer group is MISSING from this description — inevitably, we'll want to know what was the last offset committed by the group for each partition and how much it is lagging behind the latest messages in the log."
6.3 Compute consumer lag — the correct way
The deprecated approach: "In the past, the only way to get this information was to parse the commit messages that the consumer groups wrote to an internal Kafka topic. While this method accomplished its intent, Kafka does not guarantee compatibility of the internal message formats, and therefore the old method is not recommended."
Map<TopicPartition, OffsetAndMetadata> offsets =
admin.listConsumerGroupOffsets(CONSUMER_GROUP) // ①
.partitionsToOffsetAndMetadata().get();
Map<TopicPartition, OffsetSpec> requestLatestOffsets = new HashMap<>();
for (TopicPartition tp: offsets.keySet()) {
requestLatestOffsets.put(tp, OffsetSpec.latest()); // ②
}
Map<TopicPartition, ListOffsetsResult.ListOffsetsResultInfo> latestOffsets =
admin.listOffsets(requestLatestOffsets).all().get();
for (Map.Entry<TopicPartition, OffsetAndMetadata> e: offsets.entrySet()) {
String topic = e.getKey().topic();
int partition = e.getKey().partition();
long committedOffset = e.getValue().offset();
long latestOffset = latestOffsets.get(e.getKey()).offset();
System.out.println("Consumer group " + CONSUMER_GROUP
+ " has committed offset " + committedOffset
+ " to topic " + topic + " partition " + partition
+ ". The latest offset in the partition is "
+ latestOffset + " so consumer group is "
+ (latestOffset - committedOffset) + " records behind"); // ③
}① ⚠️ "unlike describeConsumerGroups, listConsumerGroupOffsets only accepts a SINGLE consumer group and not a collection."
② OffsetSpec has three very convenient implementations:
| Spec | Returns |
|---|---|
earliest() | the earliest offset in the partition |
latest() | the latest offset in the partition |
forTimestamp() | "the offset of the record written on or immediately after the time specified" |
forTimestamp() is the building block for both time-based offset resets and GDPR-style record deletion (§7.2).
6.4 Modifying consumer groups
Available operations: "deleting groups, removing members, deleting committed offsets, and modifying offsets. These are commonly used by SREs to build ad hoc tooling to recover from an emergency."
Why modify rather than delete offsets
"From all those, modifying offsets is the most useful. Deleting offsets might seem like a simple way to get a consumer to 'start from scratch,' but this really depends on the configuration of the consumer — if the consumer starts and no offsets are found, will it start from the beginning? Or jump to the latest message? Unless we have the value of
auto.offset.reset, we can't know. Explicitly modifying the committed offsets to the earliest available offsets will force the consumer to start processing from the beginning of the topic, and essentially cause the consumer to 'reset.'"
| Operation | What the consumer then does |
|---|---|
| Delete offsets | Behavior depends on the consumer's auto.offset.reset (which you may not control) → nondeterministic from the operator's view. |
| Set offsets to earliest | Deterministic. The consumer will start there. ← prefer this |
⚠️ Constraint 1: you must stop the group first
"consumer groups don't receive updates when offsets change in the offset topic. They only read offsets when a consumer is assigned a new partition or on startup. To prevent you from making changes to offsets that the consumers will not know about (and will therefore override), Kafka will PREVENT you from modifying offsets while the consumer group is active."
If the group is live, your offset write would be immediately overwritten by the running consumers' next commit.
Kafka blocks it instead of letting you shoot yourself. You get UnknownMemberIdException.
And note: "this has to be done by shutting down the consuming application directly; there is no admin command for shutting down a consumer group."
⚠️ Constraint 2: resetting offsets corrupts stateful applications
"if the consumer application maintains state (and most stream processing applications maintain state), resetting the offsets and causing the consumer group to start from the beginning can have a strange impact on the stored state."
The shoe-counting example — internalize this one:
A stream app continuously counts shoes sold. At 8:00 a.m. you find an input error and want to recalculate from 3:00 a.m.
You reset offsets to 3:00 a.m. without modifying the stored aggregate:
- stored count already includes 3:00–8:00 sales
- + reprocessing 3:00–8:00 adds them again
► You count every shoe sold today twice.
"You need to take care to update the stored state accordingly. In a development environment, we usually delete the state store completely before resetting the offsets to the start of the input topic."
Resetting offsets is not a complete reset. Offsets are only half of a stateful app's position. The other half is its state store. Reset both, or neither.
The reset code
Map<TopicPartition, ListOffsetsResult.ListOffsetsResultInfo> earliestOffsets =
admin.listOffsets(requestEarliestOffsets).all().get(); // ①
Map<TopicPartition, OffsetAndMetadata> resetOffsets = new HashMap<>();
for (Map.Entry<TopicPartition, ListOffsetsResult.ListOffsetsResultInfo> e:
earliestOffsets.entrySet()) {
resetOffsets.put(e.getKey(), new OffsetAndMetadata(e.getValue().offset())); // ②
}
try {
admin.alterConsumerGroupOffsets(CONSUMER_GROUP, resetOffsets).all().get(); // ③
} catch (ExecutionException e) {
System.out.println("Failed to update the offsets committed by group "
+ CONSUMER_GROUP + " with error " + e.getMessage());
if (e.getCause() instanceof UnknownMemberIdException) // ④
System.out.println("Check if consumer group is still active.");
}② "we convert the map with ListOffsetsResultInfo values returned by listOffsets into a map with OffsetAndMetadata values required by alterConsumerGroupOffsets." (A small type-shuffling annoyance worth knowing about.)
④ The diagnostic to remember:
"One of the most common reasons that
alterConsumerGroupOffsetsfails is that we didn't stop the consumer group first. ... If the group is still active, our attempt to modify the offsets will appear to the consumer coordinator as if a client that is not a member of the group is committing an offset for that group. In this case, we'll getUnknownMemberIdException."
UnknownMemberIdException from alterConsumerGroupOffsets = "the group is still running." That mapping is not obvious from the exception name; memorize it.
5.5 Configuration management
That parenthetical is the interesting bit: check more often than retention, because if the topic silently reverted to delete-retention, you want to notice before your data ages ou…
5.7 Cluster metadata and advanced operations
Note the honest distinction: records become inaccessible to consumers immediately, but disk cleanup is asynchronous.