--:--
notes/commonplace/kafka/kafka-producer-internals.mdx

NOTES / Commonplace / Kafka ·

Inside the Kafka producer: batching, buffering and timeouts

producer.send(record) returns at once. The record has not gone anywhere yet; it has only joined a queue inside the client. Knowing that queue explains most surprising producer behaviour.

The path of one record

  1. Serialize the key and value to bytes.
  2. Pick a partition. A key decides it: murmur2(key) % partitions (see topics, partitions and offsets). Without a key, the built-in partitioner sticks to one partition until a batch fills, and only picks among partitions that currently have a leader.
  3. Append to a batch in the record accumulator, one open batch per partition.
  4. A background sender thread drains ready batches, groups them by the broker that leads each partition, and sends one produce request per broker.
  5. The leader appends, followers fetch, and the response comes back according to acks (see replication and the ISR). The future or callback completes only now.

A batch is ready when it is full (batch.size, 16 KB) or has waited linger.ms (5 ms by default since Kafka 4.0, 0 before). A little lingering means fewer, fuller requests.

When there is no leader

The sender can only send a batch to the partition's current leader. If the partition has none (its replicas are all down, or an election is in progress), or no broker answers at all, nothing goes on the wire:

  • The batches simply stay in the accumulator. The client keeps refreshing metadata and sends them as soon as a leader appears.
  • A keyed record can't be rerouted: its key fixes the partition, and sending it elsewhere would break per-key ordering. Records without a key go to partitions that do have a leader.
  • If the producer has never seen the topic's metadata, send() itself blocks for up to max.block.ms (60 s) and then throws.
  • The accumulator has a size limit, buffer.memory (32 MB). When it is full, send() blocks for up to max.block.ms, then throws. That is how a dead cluster pushes back on the application.

The timeouts

SettingDefaultWhat it bounds
linger.ms5 msHow long a batch waits for more records before it is ready
request.timeout.ms30 sHow long one produce request waits for a response before it is retried
delivery.timeout.ms120 sThe whole life of a record after send(): waiting, sending and every retry. Then the callback gets a timeout error
max.block.ms60 sHow long send() itself may block, for metadata or for buffer space
retrieseffectively ∞Retries are bounded by delivery.timeout.ms, not by a count

delivery.timeout.ms must be at least linger.ms + request.timeout.ms. The site's simulator uses 30 s so the failure is visible without a long wait.

Retries don't create duplicates or reorder records because the producer is idempotent by default (since 3.0): each batch carries a producer id and sequence number, and the leader drops a batch it has already written.

What this means for an application

  • A successful send() call means "queued", not "stored". Check the callback or the future.
  • During an outage, records pile up in memory and the application slows down (blocked send()) before anything fails. Size buffer.memory and delivery.timeout.ms for the outage you want to ride out.
  • If a record must not be lost, use acks=all with min.insync.replicas ≥ 2, and treat a timeout as "unknown, maybe written" rather than "not written".

Try it

  • The life of one record: the happy path from serializer to commit
  • In the lab, power off brokers until a partition has no leader and watch the producer card count buffered records and time left

Related: the day it came up, Kafka day.

Check yourself1 / 3
Every broker is down. What does a producer put on the network?