5.4 Essential topic management
① describeTopics() with a list of names → DescribeTopicsResult, "which wraps a map of topic names to Future descriptions."
4.1 List topics
ListTopicsResult topics = admin.listTopics();
topics.names().get().forEach(System.out::println);
admin.listTopics()returnsListTopicsResult— "a thin wrapper over a collection of Futures."topics.names()returns a Future set of names. "When we callget()on this Future, the executing thread will wait until the server responds with a set of topic names, or we get a timeout exception."
4.2 Check existence + validate + create — the real pattern
Why not just list and search?
"One way to check if a specific topic exists is to get a list of all topics and check if the topic you need is in the list. On a large cluster, this can be inefficient. In addition, sometimes you want to check for more than just whether the topic exists — you want to make sure the topic has the right number of partitions and replicas."
The real-world example — and it's instructive because each requirement has a reason:
"Kafka Connect and Confluent Schema Registry use a Kafka topic to store configuration. When they start up, they check:
- the configuration topic exists,
- that it has only one partition — to guarantee that configuration changes will arrive in strict order,
- that it has three replicas — to guarantee availability,
- and that the topic is compacted — so the old configuration will be retained indefinitely."
A config topic's requirements, derived from first principles:
| Requirement | Why |
|---|---|
| 1 partition | Ordering is guaranteed only within a partition (Ch. 1) — config changes must be strictly ordered, so one partition. |
| 3 replicas | Availability. |
| compacted | “Retain only the last message per key” (Ch. 1) — config survives indefinitely, unlike delete-retention. |
DescribeTopicsResult demoTopic = admin.describeTopics(TOPIC_LIST); // ①
try {
topicDescription = demoTopic.values().get(TOPIC_NAME).get(); // ②
System.out.println("Description of demo topic:" + topicDescription); // ③
if (topicDescription.partitions().size() != NUM_PARTITIONS) {
System.out.println("Topic has wrong number of partitions. Exiting.");
System.exit(-1);
}
} catch (ExecutionException e) { // ④
// exit early for almost all exceptions
if (! (e.getCause() instanceof UnknownTopicOrPartitionException)) {
e.printStackTrace();
throw e;
}
// if we are here, topic doesn't exist
System.out.println("Topic " + TOPIC_NAME +
" does not exist. Going to create it now");
// Note that number of partitions and replicas is optional. If they are
// not specified, the defaults configured on the Kafka brokers will be used
CreateTopicsResult newTopic = admin.createTopics(Collections.singletonList(
new NewTopic(TOPIC_NAME, NUM_PARTITIONS, REP_FACTOR))); // ⑤
// Check that the topic was created correctly:
if (newTopic.numPartitions(TOPIC_NAME).get() != NUM_PARTITIONS) { // ⑥
System.out.println("Topic has wrong number of partitions.");
System.exit(-1);
}
}① describeTopics() with a list of names → DescribeTopicsResult, "which wraps a map of topic names to Future descriptions."
② get() on the Future gives a TopicDescription… or throws.
⚠️ ③④ The ExecutionException rule — this bites everyone once
"if the topic does not exist, the server can't respond with its description. In this case, the server will send back an error, and the Future will complete by throwing an
ExecutionException. The actual error sent by the server will be the CAUSE of the exception.""Note that ALL AdminClient result objects throw
ExecutionExceptionwhen Kafka responds with an error. This is because AdminClient results are wrappedFutureobjects, and those wrap exceptions. YOU ALWAYS NEED TO EXAMINE THECAUSEofExecutionExceptionto get the error that Kafka returned."
catch (ExecutionException e) {
Throwable actual = e.getCause(); // ← THE REAL ERROR IS HERE
if (actual instanceof UnknownTopicOrPartitionException) { ... }
}What a TopicDescription contains: "a list of all the partitions of the topic, and for each partition, in which a broker is the leader, a list of replicas and a list of in-sync replicas. Note that this does NOT include the configuration of the topic" — configuration is a separate API (§5).
⑤ Creating: "you can specify just the name and use default values for all the details. You can also specify the number of partitions, number of replicas, and the configuration."
⑥ Validating the creation: "Checking the result is more common if you relied on broker defaults when creating the topic."
⚠️ "since we are again calling
get()... this method could throw an exception.TopicExistsExceptionis common in this scenario, and you'll want to handle it (perhaps by describing the topic to check for the correct configuration)."
TopicExistsException here is the classic race: two instances of your app start simultaneously, both describe (not found), both create, one wins.
4.3 Delete topics
admin.deleteTopics(TOPIC_LIST).all().get();
// Check that it is gone. Note that due to the async nature of deletes,
// it is possible that at this point the topic still exists
try {
topicDescription = demoTopic.values().get(TOPIC_NAME).get();
System.out.println("Topic " + TOPIC_NAME + " is still around");
} catch (ExecutionException e) {
System.out.println("Topic " + TOPIC_NAME + " is gone");
}⚠️ WARNING — deletion is final
"Although the code is simple, please remember that in Kafka, deletion of topics is FINAL — there is no recycle bin or trash can to help you rescue the deleted topic, and no checks to validate that the topic is empty and that you really meant to delete it. Deleting the wrong topic could mean unrecoverable loss of data, so handle this method with extra care."
Cross-reference Ch. 2: delete.topic.enable=false is the broker-side guardrail against exactly this. If you expose deleteTopics in any tooling, put a confirmation and an audit log in front of it.
4.4 Non-blocking AdminClient — KafkaFuture.whenComplete()
When blocking get() is wrong:
"Most of the time, [blocking] is all you need — admin operations are rare, and waiting until the operation succeeds or times out is usually acceptable. There is one exception: if you are writing to a server that is expected to process a large number of admin requests. In this case, you don't want to block the server threads while waiting for Kafka to respond. You want to continue accepting requests from your users and sending them to Kafka, and when Kafka responds, send the response to the client."
vertx.createHttpServer().requestHandler(request -> { // ① Vert.x
String topic = request.getParam("topic"); // ②
String timeout = request.getParam("timeout");
int timeoutMs = NumberUtils.toInt(timeout, 1000);
DescribeTopicsResult demoTopic = admin.describeTopics( // ③
Collections.singletonList(topic),
new DescribeTopicsOptions().timeoutMs(timeoutMs));
demoTopic.values().get(topic).whenComplete( // ④ NOT get()
new KafkaFuture.BiConsumer<TopicDescription, Throwable>() {
@Override
public void accept(final TopicDescription topicDescription,
final Throwable throwable) {
if (throwable != null) {
request.response().end("Error trying to describe topic " // ⑤
+ topic + " due to " + throwable.getMessage());
} else {
request.response().end(topicDescription.toString()); // ⑥
}
}
});
}).listen(8080);"The key here is that we are not waiting for a response from Kafka.
DescribeTopicResultwill send the response to the HTTP client when a response arrives from Kafka. Meanwhile, the HTTP server can continue processing other requests."
A clever way to prove it works:
"You can check this behavior by using
SIGSTOPto pause Kafka (don't try this in production!) and send two HTTP requests to Vert.x: one with a long timeout value and one with a short value. Even though you sent the second request after the first, it will respond earlier thanks to the lower timeout value, and not block behind the first request."
5.3 Lifecycle: create, configure, close
Introduced in Kafka 2.1.0. Default behavior: "Kafka validates, resolves, and creates connections based on the hostname provided in the bootstrap server configuration (and later in…
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…