6.6 Physical storage
This single sentence explains why partition count is a capacity decision (Ch. 2's "≤6 GB per day of retention" heuristic) and why tiered storage (§6.2) is such a big deal.
6.1 The fundamental storage constraint
"The basic storage unit of Kafka is a partition replica. Partitions cannot be split between multiple brokers, and not even between multiple disks on the same broker. So the size of a partition is limited by the space available on a single mount point."
A mount point can be:
- a single disk — if JBOD configuration is used
- multiple disks — if RAID is configured
This single sentence explains why partition count is a capacity decision (Ch. 2's "≤6 GB per day of retention" heuristic) and why tiered storage (§6.2) is such a big deal.
Also:
log.dirsis "not to be confused with the location in which Kafka stores its error log, which is configured in thelog4j.propertiesfile." "The usual configuration includes a directory for each mount point that Kafka will use."
6.2 Tiered storage (planned for 3.0)
Work started late 2018. Three motivating concerns:
- “You are limited in how much data you can store in a partition. As a result, Maximum retention and partition counts aren't simply driven by product requirements but also by the limits on physical disk sizes.”
- “Your choice of disk and cluster size is driven by Storage requirements. Clusters often end up Larger than they would if latency and throughput were the main considerations, Which drives up costs.”
- “The time it takes to move partitions from one broker to another… is driven by The size of the partitions. large partitions make the cluster less elastic. These days, architectures are designed toward maximum elasticity, taking advantage of flexible cloud deployment.”
The two-tier design:
- Local serves tail reads from latency-sensitive applications → benefits from “the existing Kafka mechanism of efficiently using the page cache.”
- Remote serves backfill, and applications recovering from failure that need data older than the local tier.
- Separate retention policy per tier; local storage is “typically far more expensive than the remote tier”.
What it buys:
"allows scaling storage independent of memory and CPUs... enables Kafka to be a long-term storage solution... reduces the amount of data stored locally, and hence the amount of data that needs to be copied during recovery and rebalancing... increasing the retention period no longer requires scaling the Kafka cluster storage and adding new nodes... eliminating the need for separate data pipelines to copy the data from Kafka to external stores, as done currently in many deployments."
💡 The counterintuitive performance result (from KIP-405)
| Workload | Without tiered storage | With tiered storage |
|---|---|---|
| Case 1 — Kafka's usual high-throughput workload | p99 = 21 ms | p99 = 25 ms — slight increase, “since brokers also have to ship segments to remote storage” |
| Case 2 — some consumers are reading old data | 21 ms → 60 ms p99 — 3× degradation | 25 ms → 42 ms p99 — much less impact |
Why: “tiered storage reads are read from HDFS or S3 via a network path. Network reads do not compete with local reads on disk I/O or page cache, and leave the page cache intact with fresh data.”
"in addition to infinite storage, lower costs, and elasticity, tiered storage also delivers ISOLATION BETWEEN HISTORICAL READS AND REAL-TIME READS."
This is the real prize. The classic Kafka production incident is "someone started a backfill and now everyone's latency is terrible" — because the historical reader evicts the hot page cache and competes for disk I/O. Tiered storage moves that traffic to a different resource entirely.
Design detail: KIP-405 introduces a new component, the RemoteLogManager, and specifies its interactions with replicas catching up to the leader, and leader elections.
6.3 Partition allocation
Example: 6 brokers, 10 partitions, RF 3 → 30 replicas to place.
Three goals:
- Spread replicas evenly among brokers — here, 5 replicas per broker.
- For each partition, each replica on a different broker. “If partition 0 has the leader on broker 2, we can place the followers on brokers 3 and 4, but not on 2 and not both on 3.”
- If brokers have rack information (Kafka 0.10.0+), assign the replicas for each partition to different racks if possible → “ensures that an event that causes downtime for an entire rack does not cause complete unavailability for partitions”.
The algorithm, without rack awareness:
Start with a random broker (say 4). Assign leaders round-robin:
partition 0 leader → broker 4
partition 1 leader → broker 5
partition 2 leader → broker 0 ("because we only have 6 brokers")
...Then place the replicas at increasing offsets from the leader:
p0 leader on 4 → first follower on 5, second on 0
p1 leader on 5 → first follower on 0, second on 1
With rack awareness — a rack-alternating broker list:
So if p0's leader is on broker 2 (rack B), the first replica lands on broker 1 (rack A) — “a completely different rack. This is great, because if the first rack goes offline, we know that we still have a surviving replica, and therefore the partition is still available.”
Then: which directory?
"We do this independently for each partition, and the rule is very simple: we count the number of partitions on each directory and add the new partition to the directory with the FEWEST PARTITIONS. This means that if you add a new disk, ALL the new partitions will be created on that disk — because, until things balance out, the new disk will always have the fewest partitions."
⚠️ MIND THE DISK SPACE
"the allocation of partitions to brokers does NOT take available space or existing load into account, and the allocation of partitions to disks takes the NUMBER of partitions into account but NOT THE SIZE of the partitions. This means that if:
- some brokers have more disk space than others (perhaps because the cluster is a mix of older and newer servers),
- some partitions are abnormally large, or
- you have disks of different sizes on the same broker,
you need to be careful with the partition allocation."
This is the internals-level explanation of Ch. 2's "one disk fills while others are empty" failure. Kafka's placement is count-based and space-blind, at both levels.
6.4 File management — segments
"Because finding the messages that need purging in a large file and then deleting a portion of the file is both time-consuming and error prone, we instead split each partition into segments."
Default segment: 1 GB of data OR a week of data, whichever is smaller.
00000000000000000000.log ┐
00000000000000000000.index │ CLOSED segments —
00000000000000000000.timeindex │ eligible for deletion
00000000000001048576.log │
... ┘
00000000000009437184.log ← ACTIVE SEGMENTThe last one is the segment currently being written. “The active segment is never deleted.”
⚠️ The active-segment retention trap (restated with numbers)
"The active segment is never deleted, so if you set log retention to only store a day of data, but each segment contains five days of data, you will really keep data for FIVE DAYS because we can't delete the data before the segment is closed."
"If you choose to store data for a week and roll a new segment every day, you will see that every day we roll a new segment while deleting the oldest — so most of the time the partition will have seven segments."
That second sentence is the design target: retention period ÷ segment roll interval ≈ number of segments. Aim for a handful of segments per retention window, not one giant one.
And the file-handle consequence (from Ch. 2):
"a Kafka broker will keep an open file handle to EVERY segment in EVERY partition — even inactive segments. This leads to an unusually high number of open file handles, and the OS must be tuned accordingly."
open handles ≈ Σ over partitions of (partition_size / segment_size) × 3 — the .log/.index/.timeindex triple.
⇒ smaller segments = better retention precision But more file handles and less efficient disk writes. That's the trade.
6.5 File format — and why it's identical to the wire format
"Each segment is stored in a single data file. Inside the file, we store Kafka messages and their offsets. The format of the data on disk is IDENTICAL to the format of the messages that we send from the producer to the broker and later from the broker to the consumers."
What that buys — two things:
- Zero-copy when sending messages to consumers (§5.4)
- “avoid decompressing and recompressing messages that the producer already compressed”
And what it costs:
"if we decide to change the message format, BOTH the wire protocol AND the on-disk format need to change, and Kafka brokers need to know how to handle cases in which files contain messages of two formats due to upgrades."
Batching is mandatory since v0.11 / message format v2
"Starting with version 0.11 (and the v2 message format), Kafka producers ALWAYS send messages in batches. If you send a single message, the batching adds a bit of overhead. But with two messages or more per batch, the batching SAVES space, which reduces network and disk usage."
Three consequences the book draws out:
- “This is one of the reasons why Kafka performs better with
linger.ms=10— the small delay increases the chance that more messages will be sent together.” — the storage-level justification for Ch. 3'slinger.msadvice. - “Since Kafka creates a separate batch per partition, producers that write to fewer partitions will be more efficient as well.” — a real argument against over-partitioning, and in favour of the 2.4+ sticky partitioner (Ch. 3 §9.2).
- “if you are using compression on the producer (recommended!), sending larger batches means better compression both over the network and on the broker disks.”
Batch header contents (v2)
| Field | Notes |
|---|---|
| Magic number | current message format version |
| Offset of the first message + the difference to the last message's offset | "preserved even if the batch is later compacted and some messages are removed." First offset is set to 0 by the producer; the partition leader that first persists the batch replaces it with the real offset |
| Timestamp of the first message + the highest timestamp in the batch | "can be set by the broker if the timestamp type is set to append time rather than create time" |
| Size of the batch, in bytes | |
| The epoch of the leader that received the batch | "used when truncating messages after leader election" — KIP-101, KIP-279 |
| Checksum | validating the batch is not corrupted |
| 16 bits of attributes | compression type, timestamp type, whether the batch is part of a transaction or is a control batch |
| Producer ID, producer epoch, first sequence in the batch | "all used for exactly-once guarantees" (Ch. 8) |
| The set of messages |
Two things worth noticing: offsets are assigned by the leader, not the producer (which is why a producer can't know its offset until the ack returns), and the leader epoch is stored per batch specifically so post-election truncation can be done correctly.
Record contents — and why per-record overhead is tiny
| Field | Notes |
|---|---|
| Size of the record, in bytes | |
| Attributes | "currently there are no record-level attributes, so this isn't used" |
| The DIFFERENCE between this record's offset and the batch's first offset | delta encoding |
| The DIFFERENCE, in ms, between this record's timestamp and the batch's first timestamp | delta encoding |
| User payload: key, value, headers |
"there is very little overhead to each record, and most of the system information is at the BATCH level. Storing the first offset and timestamp of the batch in the header and only storing the DIFFERENCE in each record dramatically reduces the overhead of each record, making larger batches more efficient."
- Per batch: full offset, full timestamp, checksum, producer ID, epoch…
- Per record: deltas only — varint-friendly small numbers
⇒ bigger batches amortize a fixed header over more records and the deltas stay small. Compression then squeezes the rest.
Control batches
"Kafka also has control batches — indicating transactional commits, for instance. Those are handled by the consumer and NOT passed to the user application, and currently they include a version and a type indicator: 0 for an aborted transaction, 1 for a commit."
(These are why Ch. 1 warns that offsets are "not necessarily monotonically greater" — control batches consume offsets that consumers never surface.)
Inspecting segments yourself
bin/kafka-run-class.sh kafka.tools.DumpLogSegments"allows you to look at a partition segment in the filesystem and examine its contents." With
--deep-iteration"it will show you information about messages compressed inside the wrapper messages."
⚠️ Message format down-conversion — a CPU landmine
"Since Kafka supports upgrading brokers before all the clients are upgraded, it had to support any combination of versions... But there is a challenging situation when a NEW PRODUCER sends v2 messages to NEW BROKERS: the message is stored in v2 format, but an OLD CONSUMER that doesn't support v2 tries to read it. In this scenario, THE BROKER WILL NEED TO CONVERT the message from v2 to v1, so the consumer can parse it. This conversion uses FAR MORE CPU AND MEMORY than normal consumption, so it is best avoided."
- The conversion defeats zero-copy, and uses far more CPU and far more memory.
- Metrics to watch (KIP-188):
FetchMessageConversionsPerSecandMessageConversionsTimeMs. - “If your organization is still using old clients, we recommend CHECKING THE METRICS and UPGRADING THE CLIENTS as soon as possible.”
This is a genuinely nasty one: your brokers get slower and hotter, and the cause is a client you forgot about. Two metrics tell you instantly.
6.6 Indexes
"Kafka allows consumers to start fetching messages from any available offset. This means that if a consumer asks for 1 MB messages starting at offset 100, the broker must be able to quickly locate the message for offset 100 (which can be in ANY of the segments for the partition)."
Two indexes per partition, both segmented:
| Index | Maps |
|---|---|
| Offset index | offset → (segment file, position within the file) |
| Timestamp index | timestamp → message offset. “used when searching for messages by timestamp. Kafka Streams uses this lookup extensively, and it is also useful in some failover scenarios.” |
(The timestamp index is what powers offsetsForTimes() from Ch. 4 and OffsetSpec.forTimestamp() from Ch. 5.)
Indexes are disposable — a genuinely useful operational fact:
"Indexes are also broken into segments, so we can delete old index entries when the messages are purged. Kafka does NOT attempt to maintain checksums of the index. If the index becomes corrupted, IT WILL GET REGENERATED from the matching log segment simply by rereading the messages and recording the offsets and locations. It is also completely safe (albeit, it can cause a lengthy recovery) for an administrator to DELETE INDEX SEGMENTS if needed — they will be regenerated automatically."
Index files are derived data. Safe to delete. Never a source of truth. If you suspect index corruption: delete and restart. Cost = recovery time (see num.recovery.threads.per.data.dir).