Topic 2.1
Producer Architecture: From send() to the Broker
In one line
send() is asynchronous: the record is serialized, assigned a partition, appended to a per-partition batch in the RecordAccumulator, and later shipped by a background sender thread in produce requests grouped per broker. Acknowledgements come back through the returned future or callback.
Think of it like this
A post room. You drop letters in the outgoing tray (send returns immediately), a clerk sorts them into bags by destination city (partition batches), and a courier leaves when a bag is full or a timer rings (batch.size, linger.ms). You find out later if delivery succeeded (callback).
Key ideas
- 01
Pipeline:
ProducerRecord→ interceptors → key/value serializers → partitioner → RecordAccumulator (a buffer of batches per partition, bounded bybuffer.memory, default 32 MB) → sender thread groups ready batches by leader broker → produce request → broker response → callback. - 02
send()blocks only for metadata (first send to a topic, up tomax.block.ms, default 60 s) or when the buffer is full.send().get()makes it synchronous, which is simple and slow; use callbacks and handle errors there. - 03
The producer is thread-safe; share one instance per application (per configuration) rather than creating one per request, which would destroy batching and open many connections.
- 04
Errors in callbacks: retriable errors (leader not available, not enough replicas, timeouts) are retried automatically within
delivery.timeout.ms; non-retriable ones (record too large, serialization failure, authorization failure) fail immediately and must be handled (log, alert, send to a fallback). - 05
Close gracefully:
producer.flush()thenclose()on shutdown so buffered records are sent; otherwise records still in the accumulator are lost when the process exits.
Code & diagrams
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092,kafka2:9092,kafka3:9092");
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
p.put(ProducerConfig.ACKS_CONFIG, "all"); // default since 3.0
p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // default since 3.0
p.put(ProducerConfig.CLIENT_ID_CONFIG, "order-service");
try (KafkaProducer<String, String> producer = new KafkaProducer<>(p)) {
producer.send(new ProducerRecord<>("commerce.orders", "order-9", json), (meta, ex) -> {
if (ex != null) {
log.error("send failed for order-9", ex); // non-retriable or delivery.timeout exceeded
metrics.counter("kafka.send.failed").increment();
}
});
producer.flush(); // before shutdown
}When it breaks
A new KafkaProducer created per HTTP request
What you see
Every request opens connections and fetches metadata; batching never happens; latency and broker connection counts explode.
Fix & prevent
One long-lived producer per application, shared across threads; close it on shutdown.
Fire-and-forget send() with no callback
What you see
Records that fail after retries (or on shutdown without flush) are lost with no log or metric.
Fix & prevent
Always attach a callback that logs and counts failures; flush and close on shutdown.
Explain it without notes
Walk through what happens between calling send() and the callback firing.
Practice
Write a producer that sends 100K records with callbacks and prints the count of successes and failures, then stop one broker in the middle.
Trade-offs
- ↔
Asynchronous sends give throughput; synchronous
get()gives simple error handling at a large latency cost.
Done when you can
I can explain the producer pipeline and handle errors correctly in callbacks.