Topic 12.1
Brokers: Request Handling, Leadership, and Balance
In one line
A broker accepts client and inter-broker requests on network threads, processes them on I/O threads (appends, fetches, metadata, group and transaction coordination), stores partitions, leads some and follows others, and hosts coordinators for groups and transactions. Balanced leadership and partition placement keep load even.
Think of it like this
A post office branch with front-desk clerks taking requests (network threads), back-office staff doing the work (I/O threads), its own sorted shelves (partitions it stores), and some shelves it's the official keeper of (partitions it leads).
Key ideas
- 01
Request path: network thread reads the request into a queue → I/O (request handler) thread processes it (append to log, read from log, update metadata) → response queued → network thread writes it. Produce requests with acks=all wait in a purgatory until replication completes; fetches wait there for
fetch.min.bytes. - 02
Request types: Produce, Fetch (consumers and followers), Metadata, OffsetCommit/OffsetFetch, JoinGroup/Heartbeat (or ConsumerGroupHeartbeat in the new protocol), transactional requests.
- 03
Coordinators: group coordinators and transaction coordinators live on brokers, chosen by hashing the group ID or transactional ID to a partition of
__consumer_offsetsor__transaction_state; the leader of that partition is the coordinator. - 04
Balance: leadership should be spread evenly (preferred leaders,
auto.leader.rebalance.enable), and partitions spread by bytes and traffic, not just counts. Imbalance shows up as one broker with far higher CPU, network or request latency. - 05
Key metrics:
RequestHandlerAvgIdlePercent,NetworkProcessorAvgIdlePercent, request latency per type (TotalTimeMssplit into queue, local, remote, response times), bytes in/out per broker, leader count per broker.
Code & diagrams
kafka.network:type=RequestMetrics,name=TotalTimeMs,request=Produce p99 = 48 ms
RequestQueueTimeMs p99 = 31 ms <- waiting for an I/O thread: handler threads saturated
LocalTimeMs p99 = 4 ms <- leader append
RemoteTimeMs p99 = 11 ms <- waiting for followers (acks=all)
ResponseQueueTimeMs p99 = 1 ms
ResponseSendTimeMs p99 = 1 ms
Diagnosis: high queue time -> raise num.io.threads or reduce request rate (bigger batches).Explain it without notes
How do you tell from request metrics where produce latency comes from?
Practice
Expose broker JMX metrics in the lab (for example with the Prometheus JMX exporter) and graph Produce TotalTimeMs components under a perf test.
Trade-offs
- ↔
More I/O and network threads help saturated brokers but add context switching; tune from idle-percentage metrics.
Done when you can
I can describe the broker request path, coordinators, and how to diagnose latency by component.