How do you manage topics, configs, and ACLs programmatically with the Java AdminClient, and how do its async results behave?
answer
- Admin.create(props), AutoCloseable
- createTopics/alterConfigs/createAcls/listConsumerGroups
- *Result -> KafkaFuture; .values() per-entity vs .all()
- incrementalAlterConfigs not alterConfigs
- validateOnly = dry run; TopicExistsException via future
basics
~10 sAdminClient (Admin.create) is the Java API behind the CLI tools. You call methods like createTopics, describeTopics, alterConfigs, createAcls, listConsumerGroups. Each returns a *Result holding KafkaFutures you complete to get values or catch exceptions.
solid answer
~40 sorg.apache.kafka.clients.admin.AdminClient (interface Admin, created via Admin.create(props) with bootstrap.servers) is the programmatic equivalent of the CLI scripts. Topic management: createTopics(List<NewTopic>), deleteTopics, listTopics, describeTopics, createPartitions. Config management: describeConfigs/incrementalAlterConfigs against ConfigResource (TOPIC or BROKER). ACLs: createAcls(List<AclBinding>), deleteAcls, describeAcls. Consumer groups: listConsumerGroups, describeConsumerGroups, listConsumerGroupOffsets, alterConsumerGroupOffsets, deleteConsumerGroups. Each call returns a *Result object exposing KafkaFuture<T> (e.g. CreateTopicsResult.values() is a per-topic Map of futures, or .all() for one combined future). Operations are async and idempotent-ish: createTopics throws TopicExistsException via the future if it exists; you can pass new CreateTopicsOptions().validateOnly(true) for a dry run. Prefer incrementalAlterConfigs over the deprecated alterConfigs to avoid clobbering unspecified keys. Always close the client (it's AutoCloseable).
go deeper
Know AdminClient is the Java API for topic/config/ACL admin, created via Admin.create.
Use createTopics/describe/alter and resolve KafkaFutures, handling TopicExistsException.
Choose incrementalAlterConfigs, use validateOnly, and handle partial success via .values() vs .all().
Design provisioning/operator tooling with proper async error handling, quotas/timeouts, and authz.
## What AdminClient is `AdminClient` (program against the `Admin` interface; instantiate with `Admin.create(Properties)`) is the **Java API that the kafka-*.sh scripts themselves use**. Anything the CLI does — create topics, alter configs, manage ACLs, inspect/seek consumer groups, query log dirs — you can do in code, which is how you build self-service tooling, operators, and provisioning automation. Minimal setup: ```java Properties p = new Properties(); p.put("bootstrap.servers", "broker:9092"); try (Admin admin = Admin.create(p)) { ... } ``` It's `AutoCloseable` — always close it (try-with-resources). ## Core operation groups - **Topics:** `createTopics(Collection<NewTopic>)`, `deleteTopics`, `listTopics`, `describeTopics`, `createPartitions` (increase count). A `NewTopic(name, partitions, replicationFactor)` can also carry `.configs(Map)` and explicit replica assignments. - **Configs:** `describeConfigs(Collection<ConfigResource>)` and **`incrementalAlterConfigs`** (set/append/subtract/delete individual keys). The older `alterConfigs` **replaces the whole config set**, silently reverting any key you didn't include — prefer the incremental API. - **ACLs:** `createAcls(Collection<AclBinding>)`, `deleteAcls(Collection<AclBindingFilter>)`, `describeAcls(AclBindingFilter)`. An `AclBinding` pairs a `ResourcePattern` (e.g. TOPIC `orders`, LITERAL/PREFIXED) with an `AccessControlEntry` (principal, host, operation READ/WRITE/..., permission ALLOW/DENY). - **Consumer groups:** `listConsumerGroups`, `describeConsumerGroups`, `listConsumerGroupOffsets`, `alterConsumerGroupOffsets` (the programmatic offset reset), `deleteConsumerGroups`. - **Cluster/other:** `describeCluster`, `describeLogDirs`, `electLeaders`, `alterPartitionReassignments`. ## The async result model (the key concept) Every method returns a **`*Result`** object. It does **not** block; it wraps **`KafkaFuture`** values: ```java CreateTopicsResult r = admin.createTopics(List.of(new NewTopic("orders", 6, (short) 3))); r.values().get("orders").get(); // KafkaFuture per topic -> blocks here // or wait for all at once: r.all().get(); ``` - `.values()` → a **Map of per-entity futures** (one per topic/group/etc.), so you can handle partial success. - `.all()` → a **single future** that completes only when every entity succeeds (fails fast on the first error). - `KafkaFuture` supports `.get()` (blocking) and chaining (`thenApply`, `whenComplete`) for non-blocking flows. Because it's async, **errors surface when you resolve the future**, wrapped in `ExecutionException` whose cause is the real Kafka exception — e.g. `TopicExistsException`, `UnknownTopicOrPartitionException`, `InvalidConfigurationException`, `PolicyViolationException`, `ClusterAuthorizationException`. ## Options objects & dry runs Most methods take an options object: `new CreateTopicsOptions().validateOnly(true)` validates without creating; `timeoutMs(...)`, `retryOnQuotaViolation(...)`, and `DescribeConfigsOptions().includeSynonyms(...)` are common. `validateOnly` is the API analogue of the CLI's dry-run. ## Idempotency & gotchas - `createTopics` is **not** idempotent: re-creating an existing topic fails the future with `TopicExistsException` — catch it rather than assuming success. - Use `incrementalAlterConfigs`, not `alterConfigs`, to avoid wiping unspecified keys. - AdminClient honors security configs (`security.protocol`, SASL, SSL) just like a producer/consumer; ACL/admin calls require the right authorizations or you get `*AuthorizationException` from the future. - It connects to **brokers** via `bootstrap.servers`; in KRaft there's no ZooKeeper path. - Calls are subject to request timeouts and broker-side throttling/quotas; design retries around the future result, not by reissuing blindly.
- Why prefer incrementalAlterConfigs over alterConfigs?alterConfigs replaces the entire config set for the resource, so any key you omit reverts to default — easy to clobber existing settings. incrementalAlterConfigs applies per-key operations (SET/APPEND/SUBTRACT/DELETE) leaving other keys untouched, which is safe for partial updates.
- How do errors surface from an AdminClient call, and how do you do a dry run?Calls are async and return a *Result wrapping KafkaFutures; errors appear when you resolve the future, as an ExecutionException whose cause is the real exception (e.g. TopicExistsException). For a dry run, pass validateOnly(true) in the options (e.g. CreateTopicsOptions/incremental alter), which validates without applying.
- What's the difference between a Result's .values() and .all()?.values() returns a map of per-entity KafkaFutures so you can handle partial success (some topics created, others failed). .all() returns a single future that completes only if every entity succeeds and fails fast on the first error.
saying these in an interview costs you the question
- Saying AdminClient calls block synchronously and return values directly (they return futures).
- Using alterConfigs for partial updates and clobbering unspecified keys.
- Assuming createTopics is idempotent / silently no-ops on an existing topic.
- Forgetting to close the AutoCloseable client.
- Thinking AdminClient talks to ZooKeeper rather than brokers via bootstrap.servers.