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.
- · The accumulator is sized by
buffer.memory; a batch leaves whenbatch.sizeis full orlinger.msexpires. - · SUCCESS →
RecordMetadatawith 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.
The sequence in words
- Create a
ProducerRecord— must include topic and value; optionally key, partition, timestamp, headers. - Producer serializes key and value objects to byte arrays so they can go over the network.
- If no partition was explicitly specified, the data goes to a partitioner, which chooses a partition, usually based on the key.
- Now the producer knows the topic and partition → it adds the record to a batch of records bound for that same topic and partition.
- A separate thread sends those batches to the appropriate brokers.
- The broker responds:
- Success →
RecordMetadatawith 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.
- Success →
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 returnedFuture.