5.7 Cluster metadata and advanced operations
Note the honest distinction: records become inaccessible to consumers immediately, but disk cleanup is asynchronous.
7.1 Cluster metadata
"It is rare that an application has to explicitly discover anything at all about the cluster to which it connected. You can produce and consume messages without ever learning how many brokers exist and which one is the controller. Kafka clients abstract away this information — clients only need to be concerned with topics and partitions."
DescribeClusterResult cluster = admin.describeCluster();
System.out.println("Connected to cluster " + cluster.clusterId().get());
System.out.println("The brokers in the cluster are:");
cluster.nodes().get().forEach(node -> System.out.println(" * " + node));
System.out.println("The controller is: " + cluster.controller().get());"Cluster identifier is a GUID and therefore is not human readable. It is still useful to check whether your client connected to the correct cluster."
(That's a real safety check: "am I about to run this destructive tool against staging or prod?" — compare the cluster ID.)
Advanced operations — the SRE toolkit
"a few methods that are rarely used, and can be risky to use, but are incredibly useful when needed. Those are mostly important for SREs during incidents — but don't wait until you are in an incident to learn how to use them. Read and practice before it is too late."
7.2 Adding partitions to a topic
Why it's usually unnecessary and always risky:
"Usually the number of partitions is set when a topic is created. And since each partition can have very high throughput, bumping against the capacity limits of a topic is rare. In addition, if messages in the topic have keys, then consumers can assume that all messages with the same key will always go to the same partition and will be processed in the same order by the same consumer. For these reasons, adding partitions to a topic is rarely needed and can be risky. You'll need to check that the operation will not break any application that consumes from the topic."
(This is the same warning as Ch. 3 §9.4 — hash(key) % N changes when N changes.)
Map<String, NewPartitions> newPartitions = new HashMap<>();
newPartitions.put(TOPIC_NAME, NewPartitions.increaseTo(NUM_PARTITIONS+2));
admin.createPartitions(newPartitions).all().get();⚠️ "When expanding topics, you need to specify the TOTAL number of partitions the topic will have AFTER the partitions are added, NOT the number of new partitions."
TIP: "you may need to describe the topic and find out how many partitions exist prior to expanding it."
⚠️ "if you try to expand multiple topics at once, it is possible that some of the topics will be successfully expanded, while others will fail."
The increaseTo naming is a mercy — but the partial-failure property means multi-topic expansion is not atomic. Write your tooling to be re-runnable.
7.3 Deleting records — the compliance tool
The compliance gap, stated plainly:
"Current privacy laws mandate specific retention policies for data. Unfortunately, while Kafka has retention policies for topics, they were not implemented in a way that guarantees legal compliance. A topic with a retention policy of 30 days can store older data if all the data fits into a single segment in each partition."
(This is exactly the Ch. 2 low-volume-topic bug: retention only applies to closed segments, so a slow topic keeps data far longer than configured. Here the book names the legal consequence.)
Map<TopicPartition, ListOffsetsResult.ListOffsetsResultInfo> olderOffsets =
admin.listOffsets(requestOlderOffsets).all().get(); // forTimestamp()
Map<TopicPartition, RecordsToDelete> recordsToDelete = new HashMap<>();
for (Map.Entry<TopicPartition, ListOffsetsResult.ListOffsetsResultInfo> e:
olderOffsets.entrySet())
recordsToDelete.put(e.getKey(),
RecordsToDelete.beforeOffset(e.getValue().offset()));
admin.deleteRecords(recordsToDelete).all().get();What deleteRecords actually does:
"will mark as deleted all the records with offsets older than those specified... and make them inaccessible by Kafka consumers. The method returns the highest deleted offsets, so we can check if the deletion indeed happened as expected. Full cleanup from disk will happen asynchronously."
The two-API composition that gives you time-based deletion:
► This is how you actually implement “delete data older than 30 days” with a guarantee, instead of trusting retention config.
Note the honest distinction: records become inaccessible to consumers immediately, but disk cleanup is asynchronous. For strict "the bytes are gone" requirements, that gap matters.
7.4 Leader election
Set<TopicPartition> electableTopics = new HashSet<>();
electableTopics.add(new TopicPartition(TOPIC_NAME, 0));
try {
admin.electLeaders(ElectionType.PREFERRED, electableTopics).all().get();
} catch (ExecutionException e) {
if (e.getCause() instanceof ElectionNotNeededException) {
System.out.println("All leaders are preferred already");
}
}"If you call the command with
nullinstead of a collection of partitions, it will trigger the election type you chose for ALL partitions.""If the cluster is in a healthy state, the command will do nothing. Preferred leader election and unclean leader election only take effect when a replica other than the preferred leader is the current leader."
Type 1: Preferred leader election — safe
"Each partition has a replica designated as the preferred leader. It is preferred because if all partitions use their preferred leader replica as the leader, the number of leaders on each broker should be balanced."
"By default, Kafka will check every five minutes if the preferred leader replica is indeed the leader, and if it isn't but it is eligible to become the leader, it will elect it. If
auto.leader.rebalance.enableis false, or if you want this to happen faster,electLeader()can trigger this process."
Type 2: Unclean leader election — DATA LOSS BY DESIGN
"If the leader replica of a partition becomes unavailable, and the other replicas are NOT eligible to become leaders (usually because they are MISSING DATA), the partition will be without a leader and therefore unavailable. One way to resolve this is to trigger unclean leader election, which means electing a replica that is otherwise ineligible to become a leader as the leader anyway. THIS WILL CAUSE DATA LOSS — all the events that were written to the old leader and were not replicated to the new leader will be LOST."
The availability-vs-durability decision, made explicit:
- Wait for the old leader to come back → the partition stays unavailable (maybe indefinitely), and there is no data loss.
- Unclean leader election → the partition becomes available immediately, at the price of permanent, silent loss of everything written to the old leader that hadn't replicated.
Async caveat for both types:
"even after it returns successfully, it takes a while until all brokers become aware of the new state, and calls to
describeTopics()can return inconsistent results. If you trigger leader election for multiple partitions, it is possible that the operation will be successful for some partitions and fail for others."
7.5 Reassigning replicas
The four reasons you'd do this:
"Maybe a broker is overloaded and you want to move some replicas. Maybe you want to add more replicas. Maybe you want to move all replicas off a broker so you can remove the machine. Or maybe a few topics are so noisy that you need to isolate them from the rest of the workload."
⚠️ "reassigning replicas from one broker to another may involve copying LARGE amounts of data. Be mindful of the available network bandwidth, and throttle replication using quotas if needed; quotas are a broker configuration, so you can describe them and update them with AdminClient."
Scenario: single broker ID 0 holds one replica of every partition; a new broker (ID 1) was added.
Map<TopicPartition, Optional<NewPartitionReassignment>> reassignment = new HashMap<>();
reassignment.put(new TopicPartition(TOPIC_NAME, 0),
Optional.of(new NewPartitionReassignment(Arrays.asList(0,1)))); // ①
reassignment.put(new TopicPartition(TOPIC_NAME, 1),
Optional.of(new NewPartitionReassignment(Arrays.asList(1)))); // ②
reassignment.put(new TopicPartition(TOPIC_NAME, 2),
Optional.of(new NewPartitionReassignment(Arrays.asList(1,0)))); // ③
reassignment.put(new TopicPartition(TOPIC_NAME, 3), Optional.empty()); // ④
admin.alterPartitionReassignments(reassignment).all().get();
System.out.println("currently reassigning: " +
admin.listPartitionReassignments().reassignments().get()); // ⑤
demoTopic = admin.describeTopics(TOPIC_LIST);
topicDescription = demoTopic.values().get(TOPIC_NAME).get();
System.out.println("Description of demo topic:" + topicDescription); // ⑥| Replica list | Effect | |
|---|---|---|
| ① | [0,1] | Added a replica on the new broker (1); leader unchanged (still 0). |
| ② | [1] | Moved the single existing replica to broker 1. “Since we have only one replica, it is also the leader.” |
| ③ | [1,0] | Added a replica and made it the preferred leader. “The next preferred leader election will switch leadership to the new replica on the new broker. The existing replica will then become a follower.” |
| ④ | empty | Cancels an ongoing reassignment and “returns the state to what it was before the reassignment operation started.” |
The key insight about the list: the FIRST element of the replica list is the preferred leader. [0,1] keeps broker 0 preferred; [1,0] makes broker 1 preferred. That single ordering detail controls leadership migration.
⑥ "remember that it can take a while until it shows consistent results" — eventual consistency again.
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.
5.8 Testing with `MockAdminClient`
① "MockAdminClient is instantiated with a list of brokers (here just one), and one broker that will be our controller.