Learn Labs
13. Monitoring Kafka

13.5 JVM and OS monitoring

Via java.lang:type=OperatingSystem:

5.1 Garbage collection

For Oracle Java 1.8 with G1 (Table 13-18):

MetricJMX MBean
Full GC cyclesjava.lang:type=GarbageCollector,name=G1 Old Generation
Young GC cyclesjava.lang:type=GarbageCollector,name=G1 Young Generation

“in the semantics of GC, ‘Old’ and ‘Full’ are The same thing.”

Two attributes per metric:

CollectionCount
“the number of GC cycles of that type Since the JVM was started”
CollectionTime
“the amount of time, in Milliseconds, spent in that type of GC cycle since the JVM was started”

“As these are Counters, they can be used to tell you an Absolute number of GC cycles and Time spent in GC per unit of time. They can also provide an Average amount of time per GC cycle, though This is less useful in normal operations.”

LastGcInfo — a composite of five fields:

*"The important value to look at is duration, as this tells you how long, in milliseconds, the LAST GC cycle took. The other values (GcThreadCount, id, startTime, endTime) are informational and not very useful.

⚠️ "you will NOT be able to see the timing of EVERY GC cycle using this attribute, as YOUNG GC cycles in particular can happen FREQUENTLY."

(Why GC matters so much here: Ch. 6 §1 — a long GC pause makes a broker's ZooKeeper ephemeral node vanish, which is indistinguishable from the broker being dead.)

5.2 Java OS monitoring — two attributes worth having

Via java.lang:type=OperatingSystem:

"this information is LIMITED and does not represent everything you need to know. The two attributes that are of use, which are DIFFICULT TO COLLECT IN THE OS, are:

MaxFileDescriptorCount
“the Maximum number of FDs that the JVM is Allowed to have open”
OpenFileDescriptorCount
“the number of FDs Currently open”

“There will be FDs open for Every log segment and Every network connection, and They can add up quickly. A problem closing network connections properly could cause the broker to rapidly exhaust the number allowed.”

5.3 OS monitoring

Five areas: "CPU usage, memory usage, disk usage, disk I/O, and network usage."

CPU: "you will want to look at the system load average at the very least." Plus the percentage breakdown:

us
user space
sy
kernel space
ni
low-priority processes
id
idle
wa
Wait (on disk)
hi
hardware interrupts
si
software interrupts
st
waiting for the Hypervisor

wa and st are the two to watch on Kafka: wa means disk is the bottleneck; st means your cloud neighbor is stealing your CPU.

💡 WHAT IS SYSTEM LOAD? — most people get this wrong

*"While many know that system load is a measure of CPU usage, MOST PEOPLE MISUNDERSTAND HOW IT IS MEASURED. The load average is A COUNT OF THE NUMBER OF PROCESSES THAT ARE RUNNABLE AND ARE WAITING FOR A PROCESSOR TO EXECUTE ON. LINUX ALSO INCLUDES THREADS THAT ARE IN AN UNINTERRUPTABLE SLEEP STATE, SUCH AS WAITING FOR THE DISK.

... In a single CPU system, a value of 1 would mean the system is 100% loaded. This means that on a multiple CPU system, THE LOAD AVERAGE NUMBER THAT INDICATES 100% IS EQUAL TO THE NUMBER OF CPUs. For example, if there are 24 processors, 100% would be a load average of 24."*

Memory: "LESS IMPORTANT to track for the broker itself, as Kafka will normally be run with a relatively SMALL JVM heap size. It will use a small amount of memory outside of the heap for compression functions, but most of the system memory will be left to be used for CACHE. All the same, you should keep track of memory utilization to make sure OTHER APPLICATIONS DO NOT INFRINGE on the broker. You will also want to make sure that SWAP MEMORY IS NOT BEING USED."

Disk — the most important:

"Disk is BY FAR the most important subsystem when it comes to Kafka. All messages are persisted to disk, so the performance of Kafka depends HEAVILY on the performance of the disks."

  • Disk space and inodes (“the file and directory metadata objects for Unix filesystems”) — “as you need to assure that you are not running out of space. This is especially true for the partitions where Kafka data is being stored.”
  • Disk I/O statistics, for at least the Kafka data disks:
    • reads and writes per second
    • average read and write queue sizes
    • average wait time
    • utilization percentage

Network:

"Keep in mind that EVERY BIT INBOUND to the Kafka broker WILL BE A NUMBER OF BITS OUTBOUND EQUAL TO THE REPLICATION FACTOR of the topics, NOT INCLUDING CONSUMERS. Depending on the number of consumers, outbound network traffic could easily be AN ORDER OF MAGNITUDE LARGER than inbound. KEEP THIS IN MIND WHEN SETTING THRESHOLDS FOR ALERTS."


On this page