6. Kafka Internals
6.9 Deploy / monitor / scale / backup — through the internals lens
Ch. 6 finally makes precise what Kafka's durability actually is:
The metrics this chapter explains
| Area | Metric | What it tells you |
|---|---|---|
| Controller | ActiveControllerCount | must be exactly 1 cluster-wide |
| controller epoch changes | frequent changes = ZooKeeper instability or GC pauses (see failure #1) | |
| Replication | UnderReplicatedPartitions | ISR shrank → election ineligibility |
IsrShrinksPerSec / IsrExpandsPerSec | flapping = lag near the threshold | |
replica.lag.time.max.ms | also bounds consumer visibility delay (§5.6) — not just durability | |
| Request pipeline — the §5.1 diagram | RequestQueueSize | network threads → I/O threads |
ResponseQueueSize | I/O threads → network threads | |
NetworkProcessorAvgIdlePct | too low → add network threads | |
RequestHandlerAvgIdlePct | too low → add I/O threads | |
PurgatorySize | acks=all waits + delayed fetches | |
| Correlation ID | ties a client timeout to a broker log line | |
| Format / compat | FetchMessageConversionsPerSec | old clients burning broker CPU (KIP-188) |
MessageConversionsTimeMs | how much CPU time that conversion costs | |
| Storage | per-mount disk usage | placement is count-based, not size |
| open file descriptors | one handle per segment, always | |
| log cleaner activity / errors | offset-map memory failures |
Tuning knobs and the mechanism each controls
| Knob | The mechanism it actually touches |
|---|---|
num.network.threads | Processor threads in §5.1 — moving requests to/from queues |
num.io.threads | Request handler threads — actual request processing |
replica.lag.time.max.ms | ISR membership and the consumer visibility delay ceiling |
min.insync.replicas | The third produce-request validation (§5.3) |
acks=all | Whether the response waits in purgatory |
linger.ms | Batch fill → per-record overhead + compression ratio (§6.5) |
log.segment.bytes / log.roll.ms | Segment count → retention precision, file handles, timestamp-seek precision |
log.cleaner.enabled + offset-map memory + thread count | Whether compaction can run at all |
min/max.compaction.lag.ms | Compaction timing guarantees (compliance) |
broker.rack | The rack-alternating broker list in partition allocation (§6.3) |
client.rack + replica.selector.class | Follower fetching, at the cost of HW-propagation delay |
metadata.max.age.ms | Client metadata cache freshness → NotLeaderForPartition frequency |
"Backup" — the internals answer
Ch. 6 finally makes precise what Kafka's durability actually is:
- Kafka does not fsync. “It relies on replication for message durability.” Your durability unit is “N brokers in N failure domains hold this in page cache,” not “it's on disk.”
- Therefore: RF + rack/AZ diversity is the durability design. Correlated power loss across all ISR members can lose acknowledged data. That is by design, not a bug.
- Indexes are derived — freely deletable, auto-regenerated. They never need backing up.
- Compacted topics are the closest thing Kafka has to a permanent store: “turn Kafka into a long-term data store.” That's what makes state recovery-by-replay viable.
- Tiered storage is the real long-term answer — and explicitly aims to eliminate “separate data pipelines to copy the data from Kafka to external stores, as done currently in many deployments.”