2.8 Production concerns
Note -Xmx6g -Xms6g — min and max heap set equal, the standard practice to avoid heap resizing pauses.
8.1 Garbage collector — use G1GC
Tuning Java GC "has always been something of an art." Thankfully this changed with Java 7 and G1GC (Garbage-First). Initially considered unstable, it saw marked improvement in JDK8 and JDK11. G1GC is now recommended as the default collector for Kafka.
Why G1GC fits Kafka:
- Automatically adjusts to different workloads
- Provides consistent pause times over the application's lifetime
- Handles large heaps with ease by segmenting the heap into smaller zones and not collecting over the entire heap in each pause
Two knobs:
| Option | Default | Meaning |
|---|---|---|
MaxGCPauseMillis | 200 ms | Preferred pause time per GC cycle — not a fixed maximum; G1GC can and will exceed it if required. G1GC schedules cycle frequency and number of zones collected so each cycle takes ~this long |
InitiatingHeapOccupancyPercent | 45 | Percentage of total heap in use before G1GC starts a collection cycle. Includes both new (Eden) and old zone usage in total |
The Kafka broker is fairly efficient in how it uses heap and creates garbage objects, so it is possible to set these options lower.
Reference tuning — for a server with 64 GB memory running Kafka in a 5 GB heap:
| Option | Reference value | Default |
|---|---|---|
MaxGCPauseMillis | 20 | 200 |
InitiatingHeapOccupancyPercent | 35 | 45 — so GC runs slightly earlier |
For a server with 64 GB memory running Kafka in a 5 GB heap. The broker is “fairly efficient in how it uses heap and creates garbage objects,” so it is possible to set these lower. Note that MaxGCPauseMillis is the preferred pause time — not a fixed maximum; G1GC can and will exceed it if required.
Kafka was originally released before G1GC was available and considered stable. Therefore Kafka defaults to concurrent mark and sweep (CMS) for compatibility with all JVMs. New best practice is to use G1GC for anything Java 1.8 and later.
Change it via environment variable:
export KAFKA_JVM_PERFORMANCE_OPTS="-server -Xmx6g -Xms6g \
-XX:MetaspaceSize=96m -XX:+UseG1GC \
-XX:MaxGCPauseMillis=20 -XX:InitiatingHeapOccupancyPercent=35 \
-XX:G1HeapRegionSize=16M -XX:MinMetaspaceFreeRatio=50 \
-XX:MaxMetaspaceFreeRatio=80 -XX:+ExplicitGCInvokesConcurrent"
/usr/local/kafka/bin/kafka-server-start.sh -daemon \
/usr/local/kafka/config/server.propertiesNote -Xmx6g -Xms6g — min and max heap set equal, the standard practice to avoid heap resizing pauses.
8.2 Datacenter layout
Why it matters: in test/dev, broker physical location doesn't matter much. In production, downtime means dollars lost — through loss of service to users or loss of telemetry on what users are doing.
If not addressed prior to deploying Kafka, expensive maintenance to move servers around may be needed. A datacenter environment that has a concept of fault zones is preferable.
broker.rack — rack awareness and its sharp edge
broker.rack=rack-a # or the cloud fault domainKafka assigns new partitions in a rack-aware manner, ensuring replicas of a single partition do not share a rack.
Three critical limitations:
- This only applies to partitions that are newly created.
- The cluster does not monitor for partitions that are no longer rack aware (for example, as a result of a partition reassignment).
- Nor does it automatically correct this situation.
Therefore: use tools that keep the cluster balanced to maintain rack awareness — Cruise Control (Appendix B).
- · “This only applies to partitions that are newly created.”
- · “The cluster does not monitor for partitions that are no longer rack aware” — for example, as a result of a partition reassignment.
- · “Nor does it automatically correct this situation.”
- · Therefore: use a tool that keeps the cluster balanced to maintain rack awareness — Cruise Control (Appendix B).
Physical best practice
Best practice: each Kafka broker in a different rack — or at the very least, not sharing single points of failure for infrastructure services such as power and network.
Concretely:
- Dual power connections to two different circuits
- Dual network switches, with a bonded interface on the servers to fail over seamlessly
Even with dual connections, there is a benefit to having brokers in completely separate racks — from time to time it may be necessary to take a rack or cabinet offline for physical maintenance (moving servers, rewiring power).
8.3 Colocating applications on ZooKeeper
What Kafka actually writes to ZooKeeper: metadata about brokers, topics, and partitions. Writes only happen on changes to consumer-group membership, or changes to the Kafka cluster itself.
This traffic is generally minimal, and it does not justify a dedicated ZooKeeper ensemble for a single Kafka cluster. Many deployments use a single ensemble for multiple Kafka clusters (with a chroot path per cluster).
Kafka, ZooKeeper, and the direction of travel
Dependency on ZooKeeper is shrinking.
- Kafka 2.8.0 introduces an early-access, completely ZooKeeper-less Kafka — but it is NOT production ready. (This is KRaft.)
- Before 0.9.0.0: consumers (as well as brokers) used ZooKeeper directly to store consumer-group composition, which topics it consumed, and to periodically commit offsets per partition (to enable failover between group members).
- With 0.9.0.0: the consumer interface changed, allowing this to be managed directly with the Kafka brokers.
- Each 2.x release removes ZooKeeper from more required paths. Administration tools now connect directly to the cluster; connecting to ZooKeeper for topic creation, dynamic config changes, etc. is deprecated.
- Command-line tools moved from
--zookeeperto--bootstrap-server.--zookeeperstill works but is deprecated and will be removed.
The consumer-offsets-in-ZooKeeper problem
Even though deprecated, consumers have a configurable choice to commit offsets to ZooKeeper or Kafka, plus a configurable commit interval.
If consumers commit offsets to ZooKeeper, ZK writes per interval = (number of consumers) × (partitions each consumes).
A reasonable interval is 1 minute — because that is the window of duplicate messages a group will re-read if a consumer fails. A longer interval may be necessary if the ensemble can’t handle the traffic, “however, it is recommended that consumers using the latest Kafka libraries use Kafka for committing offsets, removing the ZooKeeper dependency.”
These commits can be a significant amount of ZooKeeper traffic, especially with many consumers. It may be necessary to use a longer commit interval if the ensemble can't handle the traffic. However, it is recommended that consumers using the latest Kafka libraries use Kafka for committing offsets, removing the ZooKeeper dependency.
Note the tension worth internalizing: commit interval is simultaneously a load knob and a correctness knob. Longer interval = less traffic and more duplicate reprocessing on failure.
Do NOT share the ensemble with non-Kafka applications
Outside of using one ensemble for multiple Kafka clusters, it is not recommended to share the ensemble with other applications if it can be avoided.
The failure chain — memorize this one:
Outside of using one ensemble for multiple Kafka clusters, it is not recommended to share the ensemble with other applications if it can be avoided. Other applications that stress the ensemble, “either through heavy usage or improper operations, should be segregated to their own ensemble.” The last step is the nasty one: the damage surfaces long after the interruption has passed.
Other applications that can stress the ensemble, either through heavy usage or improper operations, should be segregated to their own ensemble.
The nastiest part is the last step: ZooKeeper hiccups leave latent damage that surfaces later during an unrelated operation. This is why "ZK looked fine when I checked" is not exculpatory evidence during an incident review.