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?
answer
- serialize → partition → append → Sender drains
- two threads: caller vs Sender (network thread)
- send returns Future immediately
- RecordAccumulator = per-partition deques
- callback runs on Sender thread
basics
~20 ssend() 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 ssend() 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
Know that send() is async, returns a Future, and a background thread actually sends the data.
Be able to list the ordered stages (serialize, partition, append, drain, transmit, complete) and which thread each runs on.
Explain RecordAccumulator per-partition deques, batch.size vs linger.ms, and the metadata/max.block.ms wait at the front.
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.