Learn Labs
3. Kafka Producers: Writing Messages to Kafka

3.1 How does it work internally? Producer architecture

Batch vs linger

25
records/batch
200
batches/s
8.2 ms
added p99
21k/s
ceiling
Problem

Batches ship on the linger timer at 25 of 32 records — linger.ms is binding. The producer is paying per-request overhead on half-empty batches; raising linger.ms buys throughput here, at the cost of latency.

A batch is sent when it fills or when linger.ms expires, whichever happens first — so at low producer rates linger.ms sets your latency, and at high rates batch.size does.
sender threadProducerRecordtopic + value required · key, partition, timestamp, headers optionalINTERCEPTORSonSend() — before serializationSERIALIZERkey → bytes · value → bytesPARTITIONERskipped if the partition is explicitRECORD ACCUMULATORtopicA-p0 [batch][batch] · topicA-p1 [batch] · topicB-p3 [batch][batch]Broker(s)(leaders)broker responsesuccess or failureINTERCEPTORSonAcknowledgementuser Callbackruns on the producer’s main thread!APPLICATION THREADSENDER THREAD (separate)
  • · The accumulator is sized by buffer.memory; a batch leaves when batch.size is full or linger.ms expires.
  • · SUCCESS → RecordMetadata with the topic, partition, and the offset of the record within the partition.
  • · FAILURE → an error; on a retriable error the producer may retry a few more times before giving up and returning an error.
Figure 3.1.1How does it work internally? Producer architecture

The sequence in words

  1. Create a ProducerRecord — must include topic and value; optionally key, partition, timestamp, headers.
  2. Producer serializes key and value objects to byte arrays so they can go over the network.
  3. If no partition was explicitly specified, the data goes to a partitioner, which chooses a partition, usually based on the key.
  4. Now the producer knows the topic and partition → it adds the record to a batch of records bound for that same topic and partition.
  5. A separate thread sends those batches to the appropriate brokers.
  6. The broker responds:
    • Success → RecordMetadata with topic, partition, and the offset of the record within the partition.
    • Failure → an error. On a retriable error, the producer may retry a few more times before giving up and returning an error.

Two facts that surprise people:

  • A single producer object is thread-safe and can be used by multiple threads. (All chapter examples are single-threaded, but this is explicit.)
  • send() is always asynchronous internally. "Synchronous send" is just you blocking on the returned Future.

On this page