Command Palette

Search for a command to run...

Hectal
PHASE 2Intermediate ~8 min· topic 1 of 4

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.

0/4 · 0%

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

  1. 01

    Pipeline: ProducerRecord → interceptors → key/value serializers → partitioner → RecordAccumulator (a buffer of batches per partition, bounded by buffer.memory, default 32 MB) → sender thread groups ready batches by leader broker → produce request → broker response → callback.

  2. 02

    send() blocks only for metadata (first send to a topic, up to max.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.

  3. 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.

  4. 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).

  5. 05

    Close gracefully: producer.flush() then close() on shutdown so buffered records are sent; otherwise records still in the accumulator are lost when the process exits.

Code & diagrams

producer-pipeline.mermaiddiagram
Rendering diagram…
BasicProducer.javajava
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

01

Walk through what happens between calling send() and the callback firing.

Practice

01

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.