5. Managing Apache Kafka Programmatically (AdminClient)
5.9 What actually breaks in production — Ch. 5 consolidated
| # | Symptom | Root cause | Fix |
|---|---|---|---|
| 1 | createTopics().get() succeeds, then listTopics() doesn't show it | Eventual consistency — the Future completes when the controller is updated; the read went to a broker that hasn't got the metadata yet | Retry with backoff; never assert immediately after a mutation |
| 2 | deleteTopics().get() succeeds, topic still described | Same — "due to the async nature of deletes, it is possible that at this point the topic still exists" | Same |
| 3 | Catch block never matches; you can't identify the error | You checked the ExecutionException type instead of e.getCause() | Always inspect getCause() — all AdminClient results wrap errors this way |
| 4 | TopicExistsException on startup, intermittently | Race: two app instances both describe (not found), both create | Catch it and treat as success — then describe to validate config |
| 5 | Irrecoverable data loss from a wrong topic name | "deletion of topics is FINAL — no recycle bin, no checks that the topic is empty" | Broker-side delete.topic.enable=false; confirmation + audit in any tooling |
| 6 | HTTP/API server thread pool exhausted while Kafka is slow | Blocking .get() inside a request handler | KafkaFuture.whenComplete() + per-call Options().timeoutMs() |
| 7 | Service fails to start because Kafka is slow | Startup topic validation with the 120 s default request.timeout.ms, treated as fatal | Lower the timeout and start anyway, validating later or skipping |
| 8 | SASL authentication fails via a DNS alias | Client authenticates the alias; server principal is the real hostname → SASL treats the mismatch as a possible MITM | client.dns.lookup=resolve_canonical_bootstrap_servers_only |
| 9 | Client can't connect even though brokers are healthy (K8s/LB) | Client tries only the first resolved IP; that LB IP is down | client.dns.lookup=use_all_dns_ips |
| 10 | Consumer group offsets won't update; UnknownMemberIdException | The group is still active — Kafka blocks offset edits to live groups because consumers would overwrite them | Shut the consuming application down first (no admin command exists for this) |
| 11 | Stateful stream app double-counts after an offset reset | Offsets reset, but the state store still holds the old aggregate | Reset both; in dev, delete the state store entirely first |
| 12 | Deleting offsets produced unpredictable behavior | Post-delete behavior is decided by the consumer's auto.offset.reset, which the operator may not know | Set offsets explicitly (e.g. to earliest) instead of deleting |
| 13 | listConsumerGroups throws on one bad group and returns nothing | Used .all() — throws on the first error | Use .valid() (+ .errors() to inspect), for tooling that must degrade |
| 14 | Group listing empty / describe fails | Authorization, or the group's coordinator is unavailable | Check ACLs and coordinator health |
| 15 | Lag monitoring broke after a Kafka upgrade | Custom code parsed __consumer_offsets internal messages — "Kafka does not guarantee compatibility of the internal message formats" | listConsumerGroupOffsets + listOffsets |
| 16 | createPartitions created the wrong number of partitions | The argument is the TOTAL after expansion, not the count to add | describe first, then increaseTo(current + n) |
| 17 | Multi-topic expansion partially applied | "some of the topics will be successfully expanded, while others will fail" | Make tooling idempotent and re-runnable; check each result |
| 18 | Keyed consumers break after adding partitions | hash(key) % N changes (Ch. 3 §9.4) | "check that the operation will not break any application that consumes from the topic" |
| 19 | Regulator finds 90-day-old data on a 30-day-retention topic | "retention policies were not implemented in a way that guarantees legal compliance" — a single unclosed segment retains everything | listOffsets(forTimestamp) + deleteRecords(beforeOffset); also fix segment rolling (Ch. 2) |
| 20 | Records "deleted" but disk usage unchanged | "Full cleanup from disk will happen asynchronously" — deletion first only makes records inaccessible | Expected; don't gate disk-space alarms on it |
| 21 | Cluster network saturated; replication falls behind after a reassignment | Replica reassignment copies large amounts of data with no throttle | Throttle replication using quotas (broker config, editable via AdminClient) |
| 22 | Leadership didn't move after a reassignment | The first element of the replica list is the preferred leader; you kept the old broker first | Order the list intentionally; then run preferred leader election |
| 23 | ElectionNotNeededException | Cluster healthy; the preferred leader already is the leader | Not an error — handle it as a no-op |
| 24 | Partition permanently unavailable, no eligible leader | Leader down; all other replicas are missing data so are ineligible | Either wait for the old leader, or accept unclean leader election and permanent silent data loss |
| 25 | Reassignment/election results look inconsistent right after the call | Async metadata propagation | Poll listPartitionReassignments() / re-describe over time |
| 26 | Broker config file destroyed during an upgrade, no backup | No config backup process | A surviving broker IS your backup: describeConfigs against it (the book's war story) |
| 27 | A topic silently stopped being compacted and data aged out | Config drift on a topic your app depends on | Periodically validate topic config from the app — "more frequently than the default retention period, just to be safe" |
| 28 | UnsupportedOperationException: Not implemented yet in unit tests | MockAdminClient doesn't mock everything (e.g. incrementalAlterConfigs ≤ 2.5) | Inject your own implementation (Mockito doReturn) |
| 29 | MockAdminClient not found on the test classpath | It ships in a test jar | Add <classifier>test</classifier> to the dependency |
| 30 | Admin tooling broke on a ZooKeeper-less (KRaft) cluster | Code manipulated ZooKeeper directly | "NEVER use ZooKeeper directly" — AdminClient's API survives the migration |
| 31 | All mutations fail while all reads succeed | Writes go to the controller; reads go to any (least-loaded) broker → a sick controller shows exactly this asymmetry | Check controller health (describeCluster().controller()) |
| 32 | Ran a destructive tool against the wrong cluster | No cluster identity check | Compare cluster.clusterId() (a GUID) before destructive operations |