13.4 The broker metrics reference
4.1 Active controller count
- Metric name
- Active controller count
- JMX MBean
kafka.controller:type=KafkaController,name=ActiveControllerCount- Value range
- Zero or one
"At all times, ONLY ONE broker should be the controller, and ONE broker MUST ALWAYS be the controller."
| Sum across the cluster | What it means, and what to do |
|---|---|
| == 2 (two brokers claim it) | “this means that you have a problem where A controller thread that should have exited has become stuck. This can cause problems with Not being able to execute administrative tasks, such as partition moves.” Fix: “you will need to Restart both brokers at the very least.” ⚠ “when there is an Extra controller in the cluster, There will often be problems performing a safe shutdown of a broker, and You will need to force stop the broker instead.” |
| == 0 (no broker claims it) | “the cluster will Fail to respond properly in the face of state changes, including Topic or partition creation, or Broker failures.” “you must investigate further to find out Why the controller threads are not working properly. For example, A network partition from the zookeeper cluster could result in a problem like this.” “Once that underlying problem is fixed, It is wise to restart all the brokers in the cluster In order to reset state for the controller threads.” |
4.2 Controller queue size
- Metric name
- Controller queue size
- JMX MBean
kafka.controller:type=ControllerEventManager,name=EventQueueSize- Value range
- Integer, zero or more
"indicates how many requests the controller is currently waiting to process for the brokers. ... SPIKES IN THE METRIC ARE TO BE EXPECTED, BUT if this value CONTINUOUSLY INCREASES, OR STAYS STEADY AT A HIGH VALUE AND DOES NOT DROP, IT INDICATES THAT THE CONTROLLER MAY BE STUCK."
► FIX: "you will need to MOVE THE CONTROLLER to a different broker, which requires SHUTTING DOWN the broker that is currently the controller. HOWEVER, WHEN THE CONTROLLER IS STUCK, THERE WILL OFTEN BE PROBLEMS PERFORMING A CONTROLLED SHUTDOWN OF ANY BROKER."
(Ch. 12 §9.1's /admin/controller znode deletion is the alternative that doesn't require a shutdown.)
4.3 💡 Request handler idle ratio — the load metric
- Metric name
- Request handler average idle percentage
- JMX MBean
kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent- Value range
- Float, between zero and one inclusive
Why this pool and not the network threads:
"The NETWORK threads are responsible for reading and writing data to the clients across the network. This does not require significant processing, which means that EXHAUSTION OF THE NETWORK THREADS IS LESS OF A CONCERN. The REQUEST HANDLER threads, however, are responsible for servicing the client request itself, WHICH INCLUDES READING OR WRITING THE MESSAGES TO DISK. As such, as the brokers get more heavily loaded, THERE IS A SIGNIFICANT IMPACT ON THIS THREAD POOL."
The thresholds — memorize these:
- < 10% — “usually an active performance problem”
- 10–20% — “indicate a potential problem”
- > 20% — no problem indicated
INTELLIGENT THREAD USAGE
"While it may seem like you will need hundreds of request handler threads, in reality YOU DO NOT NEED TO CONFIGURE ANY MORE THREADS THAN YOU HAVE CPUs in the broker. Apache Kafka is very smart about the way it uses the request handlers, making sure to OFFLOAD TO PURGATORY those requests that will take a long time to process. This is used, for example, when requests are being quoted or when MORE THAN ONE ACKNOWLEDGMENT of produce requests is required."
(Ch. 6 §5.1's purgatory, and Ch. 6 §5.3's acks=all path. Because slow waits go to purgatory, the handler threads are never blocked waiting — so CPU count is the right sizing.)
Two causes of high utilization besides an undersized cluster:
Not enough threads
“set the number of request handler threads Equal to the number of processors in the system (Including hyperthreaded processors)”
The threads are doing unnecessary work
“Prior to kafka 0.10, the request handler thread was responsible for Decompressing every incoming message batch, Validating the messages and Assigning offsets, and then Recompressing the message batch with offsets before writing it to disk. To make matters worse, the compression methods were all behind a synchronous lock.”
“As of version 0.10, there is a new message format that allows for Relative offsets in a message batch. This means that Newer producers will set relative offsets prior to sending the batch, which allows the broker to skip recompression.”
💡 “One of the single largest performance improvements you can make is to ensure that All producer and consumer clients Support the 0.10 message format, and to Change the message format version on the brokers to 0.10 as well. This will Greatly reduce the utilization of the request handler threads.”
(This is the flip side of Ch. 6 §6.5's down-conversion problem — the relative-offset design in the v2 format is precisely what eliminates broker-side recompression.)
4.4 The rate metrics — and how to read their seven attributes
| Metric | JMX MBean |
|---|---|
| All topics bytes in | kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec |
| All topics bytes out | kafka.server:type=BrokerTopicMetrics,name=BytesOutPerSec |
| All topics messages in | kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec |
Every rate metric has seven attributes:
Descriptive (not measurements):
EventType- “the unit of measurement for all the attributes” — e.g. “bytes”
RateUnit- “the time period for the rate” — e.g. “seconds”
Four rates, different granularities:
OneMinuteRate- “An average over the previous 1 minute” — “fluctuates Quickly, provides more of a ‘Point in time’ view. Useful for seeing short spikes in traffic.”
FiveMinuteRateFifteenMinuteRate- “provide A compromise between the two”
MeanRate- “An average Since the broker was started” — “will Not vary much at all and provides an Overall trend. Though
MeanRatehas its uses, It is probably not the metric you want to be alerted on.”
And a counter:
Count- “a Constantly increasing value… since the broker was started.” “Utilized with a metrics system that supports Counter metrics, this can give you An absolute view of the measurement Instead of an averaged rate.”
⚠️ "Make sure to use the metrics APPROPRIATELY, or you will end up with a FLAWED VIEW of the broker."
Bytes out — and why it can be 6× bytes in
"The outbound bytes rate may scale DIFFERENTLY than the inbound bytes rate, thanks to Kafka's capacity to handle multiple consumers with ease. There are many deployments of Kafka where the outbound rate can easily be SIX TIMES the inbound rate! This is why it is important to observe and trend the outbound bytes rate SEPARATELY."
⚠️ REPLICA FETCHERS INCLUDED
"The outbound bytes rate ALSO INCLUDES THE REPLICA TRAFFIC. This means that if all topics are configured with a replication factor of 2, YOU WILL SEE A BYTES OUT RATE EQUAL TO THE BYTES IN RATE WHEN THERE ARE NO CONSUMER CLIENTS. If you have ONE consumer client reading all the messages, the bytes out rate will be TWICE the bytes in rate. THIS CAN BE CONFUSING when looking at the metrics if you're not aware of what is counted."
- replica fetch traffic — (RF − 1) × bytes in
- real consumers — one consumer group
bytes_out ≈ bytes_in × ( (RF − 1) + number_of_consumer_groups )- RF=2, zero consumers →
bytes_out == bytes_in(looks wrong; isn't). - RF=2, one consumer →
bytes_out == 2 × bytes_in.
Messages in — and why there's no "messages out"
WHY NO MESSAGES OUT?
"when messages are consumed, the broker just sends the NEXT BATCH to the consumer WITHOUT EXPANDING IT to find out how many messages are inside. Therefore, THE BROKER DOESN'T REALLY KNOW HOW MANY MESSAGES WERE SENT OUT. The only metric that can be provided is the number of FETCHES per second, which is a REQUEST RATE, not a messages count."
(This is zero-copy showing up in the metrics: the broker never parses the batch on the read path, so it can't count records. Ch. 6 §5.4.)
Messages-in uses: "a growth metric as a different measure of producer traffic. It can also be used in conjunction with the bytes in rate to determine an AVERAGE MESSAGE SIZE."
4.5 Partition count and leader count
Partition count kafka.server:type=ReplicaManager,name=PartitionCount
Leader count kafka.server:type=ReplicaManager,name=LeaderCountPartition count: "the total number of partitions assigned to that broker. This includes EVERY replica the broker has, regardless of whether it is a leader or follower. Monitoring this is often more interesting in a cluster that has automatic topic creation enabled, as that can leave the creation of topics outside of the control of the person running the cluster."
Leader count — and why it deserves an alert:
"It is MUCH MORE IMPORTANT to check the leader count on a regular basis, POSSIBLY ALERTING ON IT, as it will indicate when the cluster is IMBALANCED EVEN IF THE NUMBER OF REPLICAS ARE PERFECTLY BALANCED in count and size. This is because a broker can drop leadership for many reasons, such as a ZOOKEEPER SESSION EXPIRATION, and IT WILL NOT AUTOMATICALLY TAKE LEADERSHIP BACK once it recovers (except with automatic leader rebalancing). In these cases, this metric will show fewer leaders, or often ZERO, which indicates that you need to run a preferred replica election."
💡 The derived metric worth building: "use it along with the partition count to show A PERCENTAGE of partitions that the broker is the leader for. In a well-balanced cluster using replication factor 2, all brokers should be leaders for approximately 50%. If the replication factor is 3, this percentage drops to 33%."
leadership_share = LeaderCount / PartitionCount
EXPECTED ≈ 1 / RF RF=2 → 50% RF=3 → 33%- A broker at 0% just lost all leadership and never took it back.
- A broker at 100% is about to become your “one bad egg.”
4.6 Offline partitions — the "site down" metric
- Metric name
- Offline partitions count
- JMX MBean
kafka.controller:type=KafkaController,name=OfflinePartitionsCount- Value range
- Integer, zero or greater
⚠️ "This measurement is ONLY provided by the broker that is THE CONTROLLER for the cluster (all other brokers will report 0)" — so you must aggregate with
max, notsumoravg, or scrape only the controller.
Two causes:
- “All brokers hosting replicas for this partition are down.”
- “No in-sync replica can take leadership due to message-count mismatches” (with unclean leader election disabled) — Ch. 7 §3.2.
"In a production Kafka cluster, an offline partition may be impacting the producer clients, LOSING MESSAGES or CAUSING BACK PRESSURE in the application. THIS IS MOST OFTEN A 'SITE DOWN' TYPE OF PROBLEM AND WILL NEED TO BE ADDRESSED IMMEDIATELY."
4.7 Request metrics — the per-phase latency breakdown
As of version 2.5.0, metrics exist for ~50 request types, including:
AddOffsetsToTxn,AddPartitionsToTxn,AlterConfigs,AlterPartitionReassignments,ApiVersions,ControlledShutdown,CreateAcls,CreatePartitions,CreateTopics,DeleteRecords,DeleteTopics,DescribeGroups,ElectLeaders,EndTxn,Fetch,FetchConsumer,FetchFollower,FindCoordinator,Heartbeat,InitProducerId,JoinGroup,LeaderAndIsr,ListOffsets,Metadata,OffsetCommit,OffsetFetch,Produce,SaslAuthenticate,StopReplica,SyncGroup,TxnOffsetCommit,UpdateMetadata,WriteTxnMarkers, and more.
Eight metrics per request type — seven timings plus a rate. For Fetch:
kafka.network:type=RequestMetrics,name=TotalTimeMs,request=Fetchkafka.network:type=RequestMetrics,name=RequestQueueTimeMs,request=Fetchkafka.network:type=RequestMetrics,name=LocalTimeMs,request=Fetchkafka.network:type=RequestMetrics,name=RemoteTimeMs,request=Fetchkafka.network:type=RequestMetrics,name=ThrottleTimeMs,request=Fetchkafka.network:type=RequestMetrics,name=ResponseQueueTimeMs,request=Fetchkafka.network:type=RequestMetrics,name=ResponseSendTimeMs,request=Fetchkafka.network:type=RequestMetrics,name=RequestsPerSec,request=Fetch
💡 The seven phases map exactly onto Ch. 6 §5.1's threading model — this is how you localize a latency problem:
- RequestQueueTimeMs — “time in queue after received but before processing starts.”
- LocalTimeMs — “time the partition leader spends processing, including sending it to disk (but not necessarily flushing it).”
- RemoteTimeMs — “time spent waiting for the followers before processing can complete” — i.e.
acks=all. - ThrottleTimeMs — “time the response must be held to slow the requestor down to satisfy client quota settings.”
Reading it diagnostically:
| What is high | What it points to |
|---|---|
RequestQueueTimeMs | Not enough I/O (request handler) threads, or the handlers are saturated. |
LocalTimeMs | A disk problem, or the recompression tax (§4.3). |
RemoteTimeMs | Replication is slow — the followers, or their network. |
ThrottleTimeMs | You are being quota-throttled (Ch. 3 §12). |
ResponseQueueTimeMs / ResponseSendTimeMs | Network threads saturated, or slow client links. |
Attributes per timing metric: Count, Min, Max, Mean, StdDev, and percentiles 50th, 75th, 95th, 98th, 99th, 999th.
⚠️ "The metrics are ALL CALCULATED SINCE THE BROKER WAS STARTED, so keep that in mind when looking at metrics that do not change for long periods; the LONGER your broker has been running, the MORE STABLE the numbers will be."
WHAT IS A PERCENTILE?
"A 99th percentile measurement tells us that 99% of all values in the sample group are LESS THAN the value of the metric. This means that 1% of the values are GREATER. A common pattern is to view the AVERAGE value AND the 99% or 99.9% value. In this way, you can understand how the average request performs AND what the OUTLIERS are."
What to actually collect:
Minimum — for every request type:
- The average and one higher percentile (99% or 99.9%) of
TotalTimeMs RequestsPerSec
If you can — the same for the other six timing metrics: “this will allow you to narrow down any performance problems to a specific phase of request processing.”
On alert thresholds:
*"the timing metrics can be DIFFICULT. The timing for a Fetch request can vary wildly depending on settings on the client for how long it will wait for messages, how busy the particular topic is, and the speed of the network connection.
💡 "It can be very useful, however, to develop a BASELINE value for the 99.9th percentile for at least the TOTAL TIME, ESPECIALLY FOR PRODUCE REQUESTS, and alert on this. Much like the URP metric, A SHARP INCREASE IN THE 99.9TH PERCENTILE FOR PRODUCE REQUESTS CAN ALERT YOU TO A WIDE RANGE OF PERFORMANCE PROBLEMS."
Produce, not Fetch — because Fetch latency is dominated by client config (fetch.max.wait.ms), while Produce latency reflects the broker's actual health.
4.8 Topic and partition metrics
"In larger clusters these can be numerous, and it may not be possible to collect all of them as a matter of normal operations. However, they are quite useful for debugging specific issues with a client. For example, the topic metrics can be used to identify a specific topic that is causing a large increase in traffic. It also may be important to provide these metrics so that USERS of Kafka are able to access them."
Per-topic (Table 13-16) — all kafka.server:type=BrokerTopicMetrics,...,topic=TOPICNAME:
BytesInPerSec·BytesOutPerSec·MessagesInPerSecFailedFetchRequestsPerSec·FailedProduceRequestsPerSecTotalFetchRequestsPerSec·TotalProduceRequestsPerSec
"almost certainly metrics that you will NOT want to set up monitoring and alerts for. They are useful to PROVIDE TO CLIENTS, however, so that they can evaluate and debug their own usage of Kafka."
Per-partition (Table 13-17) — kafka.log:type=Log,...,topic=TOPICNAME,partition=0:
| Metric | Use |
|---|---|
Size | "the amount of data (in bytes) currently being retained on disk for the partition." Combined → "the amount of data retained for a single topic, which can be useful in ALLOCATING COSTS for Kafka to individual clients." 💡 "A DISCREPANCY BETWEEN THE SIZE OF TWO PARTITIONS FOR THE SAME TOPIC CAN INDICATE A PROBLEM WHERE THE MESSAGES ARE NOT EVENLY DISTRIBUTED ACROSS THE KEY that is being used when producing." |
NumLogSegments | "the number of log-segment files on disk. Useful along with partition size for resource tracking." |
LogEndOffset / LogStartOffset | "the highest and lowest offsets for messages in that partition." ⚠️ "the difference between these two numbers DOES NOT NECESSARILY INDICATE THE NUMBER OF MESSAGES, as LOG COMPACTION can result in 'MISSING' OFFSETS." Use case: "a more granular mapping of timestamp to offset, allowing consumers to roll back to a specific time (though this is less important with time-based index searching, introduced in Kafka 0.10.1)." |
Partition Size skew is the hot-key detector — this is how you find Ch. 3 §9.5's "Banana" problem in production.
UNDER-REPLICATED PARTITION METRICS (per-partition)
"there is a per-partition metric... In general, this is NOT VERY USEFUL in day-to-day operations, as there are TOO MANY METRICS to gather and watch. IT IS MUCH EASIER TO MONITOR THE BROKER-WIDE URP COUNT AND THEN USE THE COMMAND-LINE TOOLS to determine the specific partitions."