skip to content

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%

answer

  1. map: TopicPartition → Deque<ProducerBatch>
  2. one deque per partition
  3. batch.size (16KB) = size limit
  4. linger.ms (0) = time limit
  5. BufferPool free-list, buffer.memory cap

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.

solid answer

~50 s

The RecordAccumulator buffers records before the Sender transmits them. Internally it holds a map from TopicPartition to a Deque<ProducerBatch>; each ProducerBatch is a chunk of memory (sized around batch.size, default 16 KB) holding multiple records for one partition. send() appends to the last batch in that partition's deque, allocating a new batch if needed (and a record larger than batch.size gets its own batch, capped by max.request.size). A batch is 'ready' to drain when it is full, or when linger.ms (default 0) has elapsed since it was created, or when the buffer is exhausted / on flush/close. linger.ms adds a small artificial delay to let more records accumulate, improving batching and compression at the cost of latency. The accumulator's total memory is bounded by buffer.memory; batches are recycled via a BufferPool free-list to avoid GC churn.

go deeper

for a junior

Know the accumulator buffers records into batches before sending, one batch area per partition.

for a middle

Explain the deque-per-partition structure and that batch.size (size) and linger.ms (time) jointly trigger a send.

for a senior

Discuss the BufferPool free-list, oversized-record handling, max.request.size, and the latency/throughput trade-off of linger.ms.

for a principal

Tune batch.size/linger.ms/compression together against partition count and buffer.memory; reason about sticky partitioning's effect on batch fullness and tail latency.

## Purpose The **RecordAccumulator** is the producer's staging area. Because `send()` is asynchronous, records aren't transmitted one-by-one; they're collected here and shipped in **batches** for efficiency (fewer requests, better compression, higher throughput). ## Internal structure - A concurrent map: `TopicPartition -> Deque<ProducerBatch>`. There is **one deque per partition** so records for different partitions never contend for the same batch. - A **`ProducerBatch`** is a contiguous block of memory (target size `batch.size`) that accumulates the serialized records for a single partition, written in Kafka's record-batch format (v2), optionally compressed (`compression.type`). - Memory for batches comes from a **`BufferPool`** whose total capacity is **`buffer.memory`** (default 32 MB). Freed batches return to the pool's free list, so steady-state allocation is cheap and avoids garbage-collector pressure. ## Append path When `send()` reaches the accumulator: it locks the target partition's deque, looks at the **last** batch, and tries to append the record. If it fits, done. If not (batch full), it allocates a new batch from the BufferPool and appends there. A record **larger than `batch.size`** gets a dedicated oversized batch (still subject to `max.request.size`, default ~1 MB). ## When does a batch get sent? The Sender repeatedly asks the accumulator which partitions are **ready**. A batch is drainable when ANY of: 1. **Full** — it reached `batch.size`. (A subsequent record that wouldn't fit also marks the current batch ready.) 2. **`linger.ms` elapsed** — time since the batch was created exceeds `linger.ms` (default **0**, i.e. send as soon as the Sender can). Raising it (e.g. 5–20 ms) deliberately waits to fill batches more. 3. **Buffer pressure** — the accumulator is out of memory and must flush to free space. 4. **Explicit `flush()` / `close()`** — forces all batches ready. 5. The batch has been sitting too long and needs to respect `delivery.timeout.ms`/retry timing. ## batch.size vs linger.ms — the trade-off - `batch.size` caps how *big* a batch can get. Bigger = better amortization, more memory per partition. - `linger.ms` caps how *long* you wait to fill it. With `linger.ms=0`, a batch is sent as soon as the Sender is free even if nearly empty — low latency, smaller batches. With `linger.ms>0`, you trade a little latency for fuller batches and better compression/throughput. - They work together: a batch sends at whichever limit is hit first. ## Edge cases - **Sticky partitioning** (KIP-480/794) keeps writing to one partition until its batch fills, then rotates — this makes batches fuller even with no key. - If a record is appended and the deque already had a ready batch, the Sender may pick it up immediately. - Under `enable.idempotence=true`, batches per partition carry sequence numbers; in-flight ordering is preserved.

  • What happens to a record larger than batch.size?
    It gets its own oversized batch rather than being split; the only hard ceiling is max.request.size (default ~1 MB). batch.size is a target/minimum for normal records, not a hard per-record limit.
  • Why might you set linger.ms above 0 in production?
    To let more records accumulate per batch, which improves compression ratio and throughput and reduces request overhead, at the cost of a small, bounded added latency per record.

saying these in an interview costs you the question

  • Saying there's one global batch for all partitions — batches are per topic-partition.
  • Claiming linger.ms default is nonzero — it defaults to 0.
  • Saying batch.size is a hard per-record max — oversized records get their own batch up to max.request.size.
  • Confusing buffer.memory (total accumulator budget) with batch.size (per-batch target).

context