Learn Labs
3. Kafka Producers: Writing Messages to Kafka

3.13 What actually breaks in production — Ch. 3 consolidated

Production failure catalog
0 rows
#SymptomRoot causeFix
1Messages silently vanish; no error, no logFire-and-forget send() — errors at/after the broker are invisibleAlways use send(record, callback)
2Producer throughput is ~1/RTT messages per secondSynchronous send().get() copied from a tutorialAsync + callback
3Producer throughput collapses; the code "looks async"Blocking operation inside the callback — callbacks run on the producer's main threadHand off to another thread inside the callback
4Producer got a success ack; message is goneacks=1 (or the pre-3.0 default) + leader crashed before replicatingacks=all and min.insync.replicas≥2 (Ch. 2/7). Remember: acks=all costs no end-to-end latency
5Total data loss under broker restart, no signalacks=0Never in a durability-sensitive path
6$100 withdrawn before it was depositedretries>0 + max.in.flight>1 → retry reorderingenable.idempotence=true (ordering up to 5 in-flight, no dupes)
7Duplicate records despite everything succeedingBroker wrote + replicated, then crashed before acking; producer retried to the new leaderenable.idempotence=true
8ConfigException on producer startup after enabling idempotencemax.in.flight>5, or retries=0, or acks≠allSatisfy all three preconditions
9Producer gives up during a broker failoverdelivery.timeout.ms shorter than actual leader-election/recovery timeMeasure recovery time, set delivery.timeout.ms above it (book's example: 30s election → 120s timeout)
10Confusing timeout exceptions; can't tell where time wentSynchronous send blocks across both intervals indistinguishablyAsync + callback; reason about max.block.ms vs delivery.timeout.ms separately
11Producer startup fails with an inconsistent-timeout exceptiondelivery.timeout.ms ≤ linger.ms + request.timeout.msKeep the inequality
12Broker rejects messages: "message too large"max.request.size > broker message.max.bytesMake them match (see Ch. 2 — also replica.fetch.max.bytes and consumer fetch size)
13Poor compression ratio despite compression enabledlinger.ms=0 → batches of 1 → nothing to compress againstRaise linger.ms; compression "is much better" with real batches
14High memory use, no latency benefitbatch.size set huge, expecting it to trigger waitingbatch.size is a ceiling, not a wait trigger. Only linger.ms makes the producer wait
15send() blocks, then throws TimeoutExceptionbuffer.memory exhausted — sending faster than the broker accepts (capacity or quota)Capacity-plan; monitor; catch it at send() (not in the Future)
16Records expire before being sentSat in a batch longer than delivery.timeout.ms while the buffer was backed upsame as above
17Producer 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 keysExpected trade-off; ensure replication/availability (Ch. 7). Null keys avoid it but lose key routing
18One partition is 10× the others; a broker fills its diskHot key (the "Banana" problem)Custom partitioner isolating hot keys, or UniformStickyPartitioner if key-routing isn't needed
19After adding partitions, a key's history splits; ordering and state breakhash(key) % N changed when N changedCreate with sufficient partitions and never add them when key partitioning matters
20Cross-team byte-format war; nobody can add a fieldHand-written custom serializer; "you need to compare arrays of raw bytes" to debug; all teams must change code simultaneouslyAvro/Protobuf/Thrift + Schema Registry
21Consumers crash after a producer schema changeIncompatible schema evolution, or the deserializer lacks the writer's schemaFollow Avro compatibility rules; Schema Registry supplies the writer schema by ID
22Record size more than doubledEmbedding the full schema in every recordSchema Registry — only the schema ID travels with the record
23KafkaAvroSerializer 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
24Producer 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"
25Security alert names an IP, nobody knows which serviceclient.id unset or meaninglessSet a descriptive client.id — it drives logs, metrics, and quotas
26Client mysteriously slow; broker healthyQuota throttling — broker delaying responses and muting the channelCheck produce-throttle-time-avg/max, fetch-throttle-time-*; review quotas
27Quota change requires a full cluster restartQuotas set statically in server.propertiesUse dynamic config (kafka-configs / AdminClient)
28Interceptor leaks threads / file handlesclose() not implementedClean up in close()
29Hot-key list change requires recompiling the partitionerValues hardcoded in partition()Pass them through configure() — the book flags its own example for this