Command Palette

Search for a command to run...

Hectal
PHASE 12Advanced ~7 min· topic 1 of 4

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.

0/4 · 0%

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

  1. 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.

  2. 02

    Request types: Produce, Fetch (consumers and followers), Metadata, OffsetCommit/OffsetFetch, JoinGroup/Heartbeat (or ConsumerGroupHeartbeat in the new protocol), transactional requests.

  3. 03

    Coordinators: group coordinators and transaction coordinators live on brokers, chosen by hashing the group ID or transactional ID to a partition of __consumer_offsets or __transaction_state; the leader of that partition is the coordinator.

  4. 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.

  5. 05

    Key metrics: RequestHandlerAvgIdlePercent, NetworkProcessorAvgIdlePercent, request latency per type (TotalTimeMs split into queue, local, remote, response times), bytes in/out per broker, leader count per broker.

Code & diagrams

request-path.mermaiddiagram
Rendering diagram…
request-latency.txttext
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

01

How do you tell from request metrics where produce latency comes from?

Practice

01

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.