Walk through the full lifecycle of a single client request inside a Kafka broker, from the TCP socket to the response being written back. Name each thread pool and queue it passes through.
answer
- Acceptor -> Processor -> requestQueue -> handler -> KafkaApis -> response queue
- num.network.threads=3, num.io.threads=8, queued.max.requests=500
- SEDA staged pipeline, RequestChannel
- Purgatory releases handler for delayed ops
- Processor affinity: same processor writes the response
basics
~20 sAn acceptor thread accepts the TCP connection and hands it to a network (processor) thread, which reads the request and puts it on a shared request queue. A request-handler (I/O) thread takes it, KafkaApis processes it, and the response goes back through the network thread to the client.
solid answer
~40 sEach listener has one Acceptor thread that does nothing but accept() new TCP connections and round-robin assign each socket to a Processor (network) thread (num.network.threads, default 3). Processors run NIO selectors: they read bytes off sockets, parse them into requests, and enqueue them onto a single shared RequestChannel.requestQueue (capacity queued.max.requests, default 500). A pool of request-handler / I/O threads (num.io.threads, default 8) — the KafkaRequestHandlerPool — dequeues requests and dispatches them through KafkaApis.handle(), which routes by ApiKey (Produce, Fetch, etc.) to the actual logic (log append, replica fetch, etc.). The handler places the Response on a per-processor response queue; the owning Processor writes it back to the originating socket. Purgatory may park requests (e.g. Fetch with min.bytes) without holding a handler thread.
go deeper
Know there are two thread pools: network threads read/write sockets, I/O threads do the work; requests pass through a queue between them.
Name acceptor + processor + request queue + handler + KafkaApis in order and the relevant config defaults.
Explain purgatory releasing handler threads, processor affinity for responses, and the backpressure role of the bounded queue.
Reason about SEDA tradeoffs, capacity planning across the two pools, and which metrics localize a bottleneck to a specific stage.
Kafka's broker uses a staged, **SEDA-style (Staged Event-Driven Architecture) request pipeline** that cleanly separates network I/O from request processing so each stage can be scaled and monitored independently. ## Accepting and reading connections **1. Acceptor thread.** A *socket* is the OS endpoint of a TCP connection. For every configured listener (e.g. `PLAINTEXT://:9092`), `SocketServer` starts exactly **one Acceptor thread**. Its only job is to call `ServerSocketChannel.accept()` in a loop and assign each newly accepted connection to one of the Processor threads in **round-robin** order. Because accepting is cheap, one acceptor per listener is enough. **2. Network (Processor) threads — `num.network.threads` (default 3).** Each Processor owns a Java NIO `Selector` and a set of client connections. A Processor: - reads bytes from ready sockets, - accumulates them until a full request frame is present, - parses the bytes into a `RequestChannel.Request`, - and **enqueues** it onto the broker-wide request queue. Processors never do the actual work (no disk writes, no log reads) — they only do non-blocking network I/O so the event loop never stalls. ## Queueing and processing requests **3. Request queue.** There is **one shared `requestQueue`** inside `RequestChannel`, a bounded `ArrayBlockingQueue`. Its capacity is **`queued.max.requests` (default 500)**. When it is full, Processors stop reading new requests, which applies backpressure all the way to clients (TCP receive buffers fill, clients block/backoff). **4. Request-handler (I/O) threads — `num.io.threads` (default 8).** The `KafkaRequestHandlerPool` runs this pool. Each handler thread blocks on `requestQueue.poll()`, takes one request, and calls **`KafkaApis.handle(request)`**. Despite the name 'I/O threads', these do the *application* work: append to the log, read from the log, handle metadata/offset requests, etc. **5. KafkaApis dispatch.** `KafkaApis.handle()` is a big switch over the request's `ApiKey` (`PRODUCE`, `FETCH`, `METADATA`, `OFFSET_COMMIT`, …) that calls the matching `handleXxxRequest` method. ## Waiting and responding **6. Purgatory (edge case).** Some requests can't complete immediately: - a Produce with `acks=all` must wait for replica acks; - a Fetch with `fetch.min.bytes` may wait for data. These are parked in a **`DelayedOperationPurgatory`** and the handler thread is **released** to process other requests. When the condition is met (or a timeout fires), the operation completes and produces a response. This is why a small handler pool can serve many slow requests. **7. Response path.** The completed `Response` is placed on the **per-Processor response queue**. The owning Processor (the same one that read the request) writes the bytes back to the originating socket via its selector. This affinity keeps a connection bound to one Processor for its lifetime. ## Why this matters The split lets you tune **network capacity** (`num.network.threads`) separately from **processing capacity** (`num.io.threads`), and the bounded queue gives explicit, observable backpressure. Key metrics: - `RequestQueueSize` - `RequestHandlerAvgIdlePercent` - `NetworkProcessorAvgIdlePercent` - and per-request time breakdown (`RequestQueueTimeMs`, `LocalTimeMs`, `ResponseQueueTimeMs`).
- Which stage does the actual disk write for a produce request happen in?In a request-handler / I/O thread (num.io.threads pool), inside KafkaApis.handleProduceRequest -> ReplicaManager.appendRecords, not in the network thread.
- Why doesn't a slow acks=all produce request tie up an I/O thread for its whole duration?Because it is parked in the DelayedProduce purgatory; the handler thread returns to the pool and only resumes work when replica acks arrive or the timeout expires.
saying these in an interview costs you the question
- Saying the network threads do the log append/read (they only do network I/O).
- Claiming there is one request queue per network thread (the request queue is shared/broker-wide; response queues are per-processor).
- Thinking a handler thread blocks for the entire duration of an acks=all produce (purgatory frees it).