3. Kafka Producers: Writing Messages to Kafka
3.13 What actually breaks in production — Ch. 3 consolidated
| # | Symptom | Root cause | Fix |
|---|---|---|---|
| 1 | Messages silently vanish; no error, no log | Fire-and-forget send() — errors at/after the broker are invisible | Always use send(record, callback) |
| 2 | Producer throughput is ~1/RTT messages per second | Synchronous send().get() copied from a tutorial | Async + callback |
| 3 | Producer throughput collapses; the code "looks async" | Blocking operation inside the callback — callbacks run on the producer's main thread | Hand off to another thread inside the callback |
| 4 | Producer got a success ack; message is gone | acks=1 (or the pre-3.0 default) + leader crashed before replicating | acks=all and min.insync.replicas≥2 (Ch. 2/7). Remember: acks=all costs no end-to-end latency |
| 5 | Total data loss under broker restart, no signal | acks=0 | Never in a durability-sensitive path |
| 6 | $100 withdrawn before it was deposited | retries>0 + max.in.flight>1 → retry reordering | enable.idempotence=true (ordering up to 5 in-flight, no dupes) |
| 7 | Duplicate records despite everything succeeding | Broker wrote + replicated, then crashed before acking; producer retried to the new leader | enable.idempotence=true |
| 8 | ConfigException on producer startup after enabling idempotence | max.in.flight>5, or retries=0, or acks≠all | Satisfy all three preconditions |
| 9 | Producer gives up during a broker failover | delivery.timeout.ms shorter than actual leader-election/recovery time | Measure recovery time, set delivery.timeout.ms above it (book's example: 30s election → 120s timeout) |
| 10 | Confusing timeout exceptions; can't tell where time went | Synchronous send blocks across both intervals indistinguishably | Async + callback; reason about max.block.ms vs delivery.timeout.ms separately |
| 11 | Producer startup fails with an inconsistent-timeout exception | delivery.timeout.ms ≤ linger.ms + request.timeout.ms | Keep the inequality |
| 12 | Broker rejects messages: "message too large" | max.request.size > broker message.max.bytes | Make them match (see Ch. 2 — also replica.fetch.max.bytes and consumer fetch size) |
| 13 | Poor compression ratio despite compression enabled | linger.ms=0 → batches of 1 → nothing to compress against | Raise linger.ms; compression "is much better" with real batches |
| 14 | High memory use, no latency benefit | batch.size set huge, expecting it to trigger waiting | batch.size is a ceiling, not a wait trigger. Only linger.ms makes the producer wait |
| 15 | send() blocks, then throws TimeoutException | buffer.memory exhausted — sending faster than the broker accepts (capacity or quota) | Capacity-plan; monitor; catch it at send() (not in the Future) |
| 16 | Records expire before being sent | Sat in a batch longer than delivery.timeout.ms while the buffer was backed up | same as above |
| 17 | Producer errors on a keyed write while the cluster is "mostly fine" | Keyed records hash over ALL partitions, not just available ones — an offline partition is fatal for its keys | Expected trade-off; ensure replication/availability (Ch. 7). Null keys avoid it but lose key routing |
| 18 | One partition is 10× the others; a broker fills its disk | Hot key (the "Banana" problem) | Custom partitioner isolating hot keys, or UniformStickyPartitioner if key-routing isn't needed |
| 19 | After adding partitions, a key's history splits; ordering and state break | hash(key) % N changed when N changed | Create with sufficient partitions and never add them when key partitioning matters |
| 20 | Cross-team byte-format war; nobody can add a field | Hand-written custom serializer; "you need to compare arrays of raw bytes" to debug; all teams must change code simultaneously | Avro/Protobuf/Thrift + Schema Registry |
| 21 | Consumers crash after a producer schema change | Incompatible schema evolution, or the deserializer lacks the writer's schema | Follow Avro compatibility rules; Schema Registry supplies the writer schema by ID |
| 22 | Record size more than doubled | Embedding the full schema in every record | Schema Registry — only the schema ID travels with the record |
| 23 | KafkaAvroSerializer fails on your class | "The Avro serializer can only serialize Avro objects, not POJO" | Generate classes (avro-tools.jar / Avro Maven plugin), or use GenericRecord + explicit schema |
| 24 | Producer can't connect even though "the cluster is up" | bootstrap.servers lists a single broker, and that broker is down | "It is recommended to include at least two" |
| 25 | Security alert names an IP, nobody knows which service | client.id unset or meaningless | Set a descriptive client.id — it drives logs, metrics, and quotas |
| 26 | Client mysteriously slow; broker healthy | Quota throttling — broker delaying responses and muting the channel | Check produce-throttle-time-avg/max, fetch-throttle-time-*; review quotas |
| 27 | Quota change requires a full cluster restart | Quotas set statically in server.properties | Use dynamic config (kafka-configs / AdminClient) |
| 28 | Interceptor leaks threads / file handles | close() not implemented | Clean up in close() |
| 29 | Hot-key list change requires recompiling the partitioner | Values hardcoded in partition() | Pass them through configure() — the book flags its own example for this |