6.5 Request processing
This is the map for every broker metric and thread-pool config you'll ever tune.
"Most of what a Kafka broker does is process requests sent to the partition leaders from clients, partition replicas, and the controller."
Protocol basics:
- Binary protocol over TCP.
- "Clients always initiate connections and send requests, and the broker processes the requests and responds."
- The ordering guarantee that makes Kafka a queue: "All requests sent to the broker from a specific client will be processed in the order in which they were received — this guarantee is what allows Kafka to behave as a message queue and provide ordering guarantees on the messages it stores."
Standard request header:
| Header field | What it is for |
|---|---|
| Request type | also called the API key |
| Request version | “so the brokers can handle clients of different versions and respond accordingly” |
| Correlation ID | uniquely identifies the request; also appears in the response and in the error logs — “the ID is used for troubleshooting” |
| Client ID | identify the application that sent the request |
(Correlation ID is the underrated one: it's how you tie a client-side timeout to a specific broker-side log line.)
5.1 The threading model — memorize this diagram
This is the map for every broker metric and thread-pool config you'll ever tune.
- Processor threads take requests from client connections and place them on the request queue.
- They pick responses up off the response queue and send them back to clients.
- I/O threads pick requests off the request queue and PROCESS them.
Why purgatory exists: "At times, responses to clients have to be delayed — consumers only receive responses when data is available, and admin clients receive a response to a DeleteTopic request after topic deletion is underway. The delayed responses are held in a purgatory until they can be completed."
(Purgatory is also where acks=all produce requests wait — §5.3.)
5.2 Request routing — metadata requests and the "Not a Leader" loop
The three most common client request types:
| Type | Sent by | Contains |
|---|---|---|
| Produce | producers | messages the clients write |
| Fetch | consumers AND follower replicas | read requests |
| Admin | admin clients | metadata operations (create/delete topics) |
The routing constraint:
"Both produce requests and fetch requests have to be sent to the LEADER replica of a partition. If a broker receives a produce request for a specific partition and the leader for this partition is on a different broker, the client will get an error response of 'Not a Leader for Partition.' The same error will occur if a fetch request... arrives at a broker that does not have the leader. Kafka's clients are responsible for sending produce and fetch requests to the broker that contains the leader for the relevant partition."
- ③ The client caches the metadata and routes produce/fetch to the right broker.
- Refresh triggers: periodically —
metadata.max.age.ms. - And immediately on a “Not a Leader” error: “the error indicates the client is using outdated information and is sending requests to the wrong broker” → refresh, then retry.
This is the machinery behind Ch. 3's "retriable error." NotLeaderForPartition is retriable precisely because the client's recovery action is refresh metadata and retry — and the underlying condition (a leader election) resolves on its own.
5.3 Produce request handling
Validations the leader runs first:
- Does the user sending the data have write privileges on the topic?
- Is the number of
acksspecified valid? (only 0, 1, and “all”) - If
acks=all: are there enough in-sync replicas for safely writing? “Brokers can be configured to refuse new messages if the number of in-sync replicas falls below a configurable number” — this ismin.insync.replicas(Ch. 2, Ch. 7).
⚠️ The durability model, stated bluntly
"Then the broker will write the new messages to local disk. On Linux, the messages are written to the FILESYSTEM CACHE, and there is NO GUARANTEE about when they will be written to disk. Kafka DOES NOT WAIT for the data to get persisted to disk — IT RELIES ON REPLICATION FOR MESSAGE DURABILITY."
Kafka's durability is not “the bytes are on the platter.” Kafka's durability is “N brokers have it in their page cache.”
- This is why replication factor is the durability knob, not
fsync. - This is why Ch. 2 says: if you raise
vm.dirty_ratio, “it is highly recommended that replication be used.” - This is why a simultaneous power loss across all ISR brokers can lose acknowledged data. Rack/AZ diversity is the answer.
Then the acks branch:
(So acks=all latency is literally "time spent in purgatory." That's the metric to watch.)
5.4 Fetch request handling
The request shape:
"something like 'Please send me messages starting at offset 53 in partition 0 of topic Test and messages starting at offset 64 in partition 3 of topic Test.'"
Why the upper limit is mandatory, not optional:
"Clients also specify a limit to how much data the broker can return for each partition. The limit is important because clients need to allocate memory that will hold the response. Without this limit, brokers could send back replies large enough to cause clients to run out of memory."
Validation: "does this offset even exist for this particular partition? If the client is asking for a message that is so old it got deleted from the partition, or an offset that does not exist yet, the broker will respond with an error."
💡 Zero-copy — why the on-disk format equals the wire format
"Kafka famously uses a zero-copy method to send the messages to the clients — this means that Kafka sends messages from the file (or more likely, the Linux filesystem cache) DIRECTLY TO THE NETWORK CHANNEL without any intermediate buffers. This is different than most databases where data is stored in a local cache before being sent to clients. This technique removes the overhead of copying bytes and managing buffers in memory, and results in much improved performance."
The precondition, from §7.2: the on-disk format is identical to the wire format. If the broker had to transform bytes (decompress/recompress, or down-convert message versions — §7.3), zero-copy is defeated. That single fact links compression choice, client version skew, and broker CPU.
5.5 Lower boundary (fetch.min.bytes) — the mechanism
"clients can also set a lower boundary on the amount of data returned. Setting the lower boundary to 10K is the client's way of telling the broker, 'Only return results once you have at least 10K bytes to send me.' This is a great way to reduce CPU and network utilization when clients are reading from topics that are not seeing much traffic."
“The same amount of data is read overall but with much less back-and-forth and therefore less overhead.”
Plus the timeout: "If you didn't satisfy the minimum amount of data to send within x milliseconds, just send what you got." (= fetch.max.wait.ms, Ch. 4.)
5.6 ⚠️ The high-water mark — consumers can't read everything on the leader
"It is interesting to note that not all the data that exists on the leader of the partition is available for clients to read. Most clients can only read messages that were written to ALL IN-SYNC REPLICAS (follower replicas, even though they are consumers, are exempt from this — otherwise replication would not work)... until a message was written to all in-sync replicas, it will not be sent to consumers — attempts to fetch those messages will result in an EMPTY RESPONSE rather than an error."
The consistency argument — this is the reasoning that matters:
"messages not replicated to enough replicas yet are considered 'unsafe' — if the leader crashes and another replica takes its place, these messages will no longer exist in Kafka. If we allowed clients to read messages that only exist on the leader, we could see inconsistent behavior. For example, if a consumer reads a message and the leader crashed and no other broker contained this message, the message is gone. No other consumer will be able to read this message, which can cause inconsistency with the consumer who did read it."
- A consumer asking for offset 6 gets an empty response — not an error.
- Why: if the leader dies now and a follower with only 0–5 takes over, message 6 never existed as far as the cluster is concerned. A consumer that had read it would hold a phantom record no other consumer could ever see.
⚠️ The latency consequence: "if replication between brokers is slow for some reason, it will take longer for new messages to arrive to consumers (since we wait for the messages to replicate first). This delay is limited to
replica.lag.time.max.ms— the amount of time a replica can be delayed while still being considered in sync."
This is the mechanism behind Ch. 3's biggest insight: end-to-end latency is identical for all acks settings, because visibility is gated on ISR replication regardless of what the producer waited for. Now you know exactly why: the high-water mark.
And it gives you a latency bound worth remembering: worst-case added consume latency from slow replication ≈ replica.lag.time.max.ms. Beyond that, the slow replica drops out of the ISR and stops holding the HW back.
5.7 Fetch session cache — incremental fetch requests
The problem:
"In some cases, a consumer consumes events from a large number of partitions. Sending the list of all the partitions it is interested in to the broker with every request and having the broker send all its metadata back can be very inefficient — the set of partitions rarely changes, their metadata rarely changes, and in many cases there isn't that much data to return."
The solution:
A consumer creates a cached session storing the list of partitions it consumes from plus their metadata. From then on:
- subsequent requests are incremental fetch requests
- consumers no longer specify all partitions each time
- “Brokers will only include metadata in the response if there were any changes”
⚠ “The session cache has limited space, and Kafka prioritizes follower replicas and consumers with a large set of partitions, so in some cases a session will not be created or will be evicted. In both cases the broker returns an appropriate error, and the consumer transparently resorts to full fetch requests.”
Note the graceful degradation: eviction is not an outage, it's a silent efficiency loss. Which means it's the kind of thing you only find via metrics.
5.8 The protocol surface and version negotiation
"The Kafka protocol currently handles 61 different request types, and more will be added. Consumers alone use 15 request types to form groups, coordinate consumption, and allow developers to manage the consumer groups."
The same protocol is used broker-to-broker: "Those requests are internal and should not be used by clients. For example, when the controller announces that a partition has a new leader, it sends a LeaderAndIsr request to the new leader (so it will know to start accepting client requests) and to the followers (so they will know to follow the new leader)."
Protocol evolution — the ZooKeeper-removal trail:
Consumer offsets moved from ZooKeeper to a special Kafka topic. That required new requests: OffsetCommitRequest, OffsetFetchRequest, ListOffsetsRequest. “Now when an application calls the client API to commit consumer offsets, the client no longer writes to ZooKeeper; instead, it sends OffsetCommitRequest to Kafka.”
Topic creation moved from the CLI writing to ZooKeeper to CreateTopicRequest. “it allows clients in languages that don't have a ZooKeeper library to create topics by asking Kafka brokers directly.”
⚠️ 5.9 Why you must upgrade brokers BEFORE clients
The book walks a concrete example (Metadata request v0 → v1, adding controller info between 0.9.0 and 0.10.0):
- The old client doesn't expect controller info “and wouldn't know how to parse it anyway.”
- “This is the reason we recommend upgrading the brokers before upgrading any of the clients — new brokers know how to handle old requests, but not vice versa.”
Mitigation added in 0.10.0: ApiVersionRequest — "allows clients to ask the broker which versions of each request are supported and to use the correct version accordingly. Clients that use this new capability correctly will be able to talk to older brokers."
Future work: KIP-584 — APIs for clients to discover which features brokers support, and for brokers to gate features per version. "at this time it seems likely to be part of version 3.0.0."
6.4 Replication
Scale note: "each broker typically stores hundreds or even thousands of replicas belonging to different topics and partitions."
6.6 Physical storage
This single sentence explains why partition count is a capacity decision (Ch. 2's "≤6 GB per day of retention" heuristic) and why tiered storage (§6.2) is such a big deal.