skip to content

Send Flow and Record Accumulator

What actually happens after send() returns: serialization, the per-partition accumulator, and the Sender thread draining batches. Interviewers ask because it explains why send() is asynchronous and when it blocks on buffer.memory.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

Walk through what happens when you call KafkaProducer.send(record). What are the stages between the call returning and the record actually reaching the broker?

level: juniorimportance: must knowfreq 78%

answer

  1. serialize → partition → append → Sender drains
  2. two threads: caller vs Sender (network thread)
  3. send returns Future immediately
  4. RecordAccumulator = per-partition deques
  5. callback runs on Sender thread

basics

~20 s

send() is asynchronous. The producer serializes the key/value, picks a partition, and appends the record to an in-memory buffer (RecordAccumulator). A background Sender thread later batches and transmits it. send() returns a Future immediately, before the broker has the record.

solid answer

~40 s

send() does NOT block on the network. On the calling thread it: (1) waits for metadata if the topic is unknown (bounded by max.block.ms), (2) serializes key and value with the configured serializers, (3) runs the Partitioner to choose a partition, (4) appends the serialized record into a per-partition batch inside the RecordAccumulator (an in-memory buffer governed by buffer.memory). It then returns a Future<RecordMetadata>. A separate background thread, the Sender (running the KafkaProducer's I/O loop over NetworkClient), drains ready batches, groups them per broker, and sends produce requests. When the broker acks, the Sender completes the Future and fires any Callback. So the work splits between the user thread (serialize/partition/append) and the Sender thread (drain/transmit/complete).

go deeper

for a junior

Know that send() is async, returns a Future, and a background thread actually sends the data.

for a middle

Be able to list the ordered stages (serialize, partition, append, drain, transmit, complete) and which thread each runs on.

for a senior

Explain RecordAccumulator per-partition deques, batch.size vs linger.ms, and the metadata/max.block.ms wait at the front.

for a principal

Reason about failure surfaces (sync serializer errors vs async retriable errors), callback-thread constraints, and how the design decouples app throughput from broker latency.

## What `send()` actually is A Kafka *producer* is a client library that ships records (key/value messages) to Kafka *brokers* (the servers). `KafkaProducer.send(record)` looks synchronous but is fundamentally **asynchronous**: it hands the record off to an in-memory buffer and returns a `Future` immediately, long before the record is durably stored on a broker. ## The two threads 1. **The application (caller) thread** — the thread that calls `send()`. 2. **The Sender thread** (sometimes called the I/O thread, named `kafka-producer-network-thread`) — a single background thread the producer starts internally. It owns the `NetworkClient` and does all socket I/O. ## Stages on the caller thread (inside `send()`) 1. **Metadata wait.** The producer needs cluster metadata (which broker leads each partition). If metadata for the topic is missing/stale, `send()` *blocks* here until metadata arrives, bounded by **`max.block.ms`** (default 60000 ms). If it times out, send throws `TimeoutException`. 2. **Serialization.** The configured `key.serializer` and `value.serializer` convert objects to `byte[]`. This runs on the caller thread, so a slow/throwing serializer slows or fails `send()` synchronously. 3. **Partitioning.** The `Partitioner` chooses a partition: if the record has an explicit partition it's used; else if it has a key, partition = hash(key) % numPartitions (murmur2); else the modern default uses *sticky* batching (`DefaultPartitioner`/`UniformStickyPartitioner` behavior, KIP-480/KIP-794) to fill a batch before rotating. 4. **Append to RecordAccumulator.** The serialized record is appended to a `ProducerBatch` in a **per-partition deque** inside the `RecordAccumulator`. This buffer's total size is capped by **`buffer.memory`** (default 32 MB). A batch fills up to **`batch.size`** (default 16 KB) or is sent earlier after **`linger.ms`**. 5. **Return a Future.** `send()` returns `Future<RecordMetadata>`. Nothing has hit the network yet. ## Stages on the Sender thread 6. **Drain.** The Sender polls the accumulator for batches that are *ready* (full, or `linger.ms` elapsed, or buffer under pressure), groups them by destination broker, and builds `ProduceRequest`s. 7. **Transmit.** `NetworkClient` writes requests over the socket; at most `max.in.flight.requests.per.connection` (default 5) requests are outstanding per connection. 8. **Completion.** When the broker responds (per `acks`), the Sender completes each record's `Future` and invokes its `Callback` — **on the Sender thread**, so callbacks must be fast and non-blocking. ## Edge cases - If the buffer is full, the *append* step (4) blocks up to `max.block.ms`, then throws. - A retriable error (e.g. `NotLeaderForPartition`) causes the Sender to retry the batch; the Future isn't completed until success or final failure. - `flush()` blocks the caller until all buffered records are sent/failed; `close()` drains then stops the Sender.

  • Which steps run on the caller's thread and which on the Sender thread?
    Serialization, partitioning, and appending to the accumulator run on the caller thread (and the metadata wait). Draining batches, network transmission via NetworkClient, and completing the Future/firing the Callback run on the Sender (background I/O) thread.
  • If your serializer throws, where and when do you see the error?
    Synchronously, on the calling thread, inside the send() call itself — serialization happens before the record is buffered, so it surfaces immediately rather than via the Future/Callback.

saying these in an interview costs you the question

  • Saying send() blocks until the broker acknowledges — it does not by default; it returns a Future.
  • Claiming the network send happens on the caller thread — transmission is on the Sender thread.
  • Thinking the Callback runs on the application thread — it runs on the Sender thread.
  • Forgetting that serialization/partitioning happen before buffering, on the caller thread.

context

open as a page

Describe the RecordAccumulator: how does it organize buffered records, and how do batch.size and linger.ms govern when a batch is sent?

level: middleimportance: must knowfreq 70%

basics

~20 s

The RecordAccumulator is the producer's in-memory buffer. It keeps one deque of ProducerBatches per topic-partition. A record is appended to the partition's current batch. A batch becomes sendable when it fills to batch.size or when linger.ms elapses, whichever comes first.

open as a page

How does buffer.memory and max.block.ms create back-pressure on a producer, and what does the application experience when the buffer fills?

level: seniorimportance: must knowfreq 66%

basics

~20 s

buffer.memory caps the total bytes the accumulator can hold. If the app produces faster than the Sender can drain, the buffer fills. Then send() blocks waiting for free space for up to max.block.ms; if it can't get memory in time, send() throws a TimeoutException.

open as a page

Compare the two ways to get the result of a send: the returned Future versus a Callback. When is each appropriate, and what are the pitfalls?

level: middleimportance: should knowfreq 58%

basics

~20 s

send() returns a Future<RecordMetadata>; calling future.get() blocks until the broker responds, turning async into sync. Alternatively pass a Callback to send(record, callback) that the producer invokes asynchronously on completion. Use Future.get() for sync confirmation, Callback for non-blocking handling.

open as a page

How does the Sender thread drain the RecordAccumulator and use the NetworkClient, and how do max.in.flight.requests.per.connection and idempotence interact with batch ordering and retries?

level: principalimportance: should knowfreq 48%

basics

~20 s

The Sender thread polls the accumulator for ready batches, groups them by leader broker, and sends ProduceRequests via the NetworkClient — up to max.in.flight.requests.per.connection outstanding per connection. With idempotence on, the producer can keep 5 in-flight and still preserve per-partition order and dedup on retries; without it, retries can reorder unless in-flight is 1.

open as a page