Most explanations of Kafka stop at a picture of a topic split into numbered boxes. That picture is right but not useful for operations. The decisions that hurt in production, such as how many partitions to create, why a broker restart takes ten minutes, why a disk fills even though retention is one day, or why adding partitions broke a downstream join, all come from what a partition is physically: a directory of append-only files, replicated to several brokers, with one leader at a time.
This article treats the partition as a storage and placement unit. It covers the files on disk and how a read finds its bytes, how retention and compaction act on segments, how the high watermark decides what consumers can see, where replicas are placed, and how to size and change partition counts. Key-to-partition mapping is covered in Kafka partitioning strategies, the in-sync replica protocol in ISR replication, and ordering guarantees in Kafka ordering architecture. Defaults quoted here are Apache Kafka's documented defaults; your distribution or managed service may override them.
Three jobs, one unit
A partition is simultaneously three things, and the tension between them explains most sizing trade-offs.
It is the unit of order. Records within a partition have consecutive offsets, and every consumer reads them in that order. There is no order across partitions.
It is the unit of parallelism. Within a consumer group, each partition is read by exactly one consumer at a time, so the partition count is the upper bound on useful consumers. Producers also spread load across partition leaders on different brokers.
It is the unit of replication and placement. Each partition has a replica set, and leadership, replication, recovery and reassignment are all done per partition. Every partition costs open files, memory for indexes, metadata in the controller, and failover work.
More partitions buy parallelism and spread, and cost per-partition overhead and wider failure blast radius. Fewer partitions are cheaper and recover faster but cap consumer throughput. There is no free choice, only a sized one.
On disk: segments and indexes
Each replica of partition 7 of topic orders is a directory named orders-7 under the broker's log.dirs. Inside, the log is split into segments. Each segment is a .log file of record batches plus an .index and a .timeindex, all named by the segment's base offset padded to 20 digits, for example 00000000000000041210.log. There is a leader-epoch-checkpoint file, producer-state .snapshot files for idempotence and transactions, and a .txnindex where aborted transactions exist.
Only the newest segment, the active segment, accepts appends. It rolls into a closed segment when it reaches segment.bytes (1 GiB by default) or segment.ms (seven days by default), whichever comes first. Closed segments are never modified, only deleted or replaced as a whole.
The offset index is sparse. Roughly every index.interval.bytes (4,096 by default) of log data, the broker writes one entry mapping a relative offset to a byte position. To serve a fetch at offset 52,000, the broker picks the segment whose base offset is the largest value not above 52,000, binary-searches its index for the closest entry at or below the target, and scans forward a few kilobytes. The time index does the same for timestamps, which is how offsetsForTimes and time-based retention work. Indexes are memory-mapped, so a partition being read costs page cache and file handles even when idle.
# inspect one segment (Kafka 3.x+/4.x tooling)
bin/kafka-dump-log.sh --files /var/kafka/orders-7/00000000000000041210.log --print-data-log | head
bin/kafka-dump-log.sh --files /var/kafka/orders-7/00000000000000041210.index
# how a fetch at a target offset finds its bytes (logic, simplified)
def locate(segments, target):
seg = max((s for s in segments if s.base_offset <= target), key=lambda s: s.base_offset)
rel = target - seg.base_offset
pos = seg.index.floor(rel).position # sparse entry at or below rel
return seg.log.scan_from(pos, until_offset=target)
Retention and compaction act on segments
Because closed segments are immutable, retention deletes whole segments. A segment is eligible when its newest record is older than retention.ms, or when the partition exceeds retention.bytes. The active segment is never deleted. A low-traffic partition with retention of one day but segment.ms of seven days can therefore keep records for up to about eight days, because its only segment stays active. Disks fill on topics that look as if they are within retention for exactly this reason.
Compaction (cleanup.policy=compact) rewrites closed segments so that only the latest record per key survives, plus tombstones for a configurable time. It also leaves the active segment alone, so recently updated keys keep several versions until the segment rolls. Compaction preserves offsets: surviving records keep their original offsets, which leaves gaps a consumer must tolerate.
Tiered storage (KIP-405, production-ready from Kafka 3.9) changes where closed segments live, not what they are. Closed segments are copied to object storage and deleted locally once a local-retention limit is reached. The partition model is unchanged; only the placement of old bytes moves.
Visibility: log end offset and high watermark
Every replica tracks its log end offset (LEO), the offset it will assign or store next. The leader also tracks the high watermark (HW): the highest offset that every in-sync replica has stored. Consumers on the default read_uncommitted isolation can read only below the HW, so a consumer never reads a record that a leader failover could erase. Transactional consumers on read_committed stop at the last stable offset, which is at or below the HW.
A producer using acks=all gets its acknowledgement once the record is below the HW. min.insync.replicas sets how small the in-sync set may shrink before writes are refused. With replication factor 3 and min ISR 2, one broker can be lost without losing acknowledged data or availability. A follower drops out of the in-sync set when it has not caught up within replica.lag.time.max.ms, which defaults to 30 seconds since Kafka 2.5.
Leader epochs fix the subtle part. Each leadership change increments an epoch, and the checkpoint file records the first offset of each epoch. A replica returning after failure asks the new leader where its epoch ended and truncates any divergent tail, instead of trusting its own high watermark. That is the mechanism behind the guarantee that acknowledged records under acks=all survive failover. The full protocol is in the ISR article linked above.
Placement: racks, leaders and follower fetching
When a topic is created, replicas are assigned across brokers. With broker.rack set, Kafka spreads each partition's replicas over different racks or availability zones, so losing one zone loses at most one replica per partition. Set it before creating topics; existing assignments do not move by themselves.
The first replica in each assignment is the preferred leader. After a broker restarts, its partitions' leadership has moved elsewhere. With auto.leader.rebalance.enable (true by default), the controller moves leadership back once imbalance passes a threshold, so produce load returns to an even spread.
Consumers normally fetch from the leader, which in a multi-zone cluster means cross-zone traffic for two out of three consumers. Since Kafka 2.4 (KIP-392), brokers configured with replica.selector.class=org.apache.kafka.common.replica.RackAwareReplicaSelector let consumers that set client.rack read from a replica in their own zone. Followers serve only up to the high watermark they have learned, so reads see slightly higher latency in exchange for lower network cost.
Worked example: sizing a topic
A team needs a topic for payment events. Peak production is 60 MB/s. Load tests show one producer instance pushes about 25 MB/s per partition leader without latency degrading, and one consumer instance, doing a database write per record, processes 4 MB/s. Order is required per account.
The classic rule is partitions at least the larger of target throughput over per-partition producer throughput, and target throughput over per-consumer throughput. That gives max(60/25, 60/4) = max(2.4, 15) = 15. Consumer processing is the binding constraint, as it usually is.
Then add headroom for growth, because changing the count later is disruptive (next section). If traffic is expected to double in two years, plan for 30 consumers. Rounding to a number with many divisors, such as 36 or 48, lets consumer groups of 6, 12 or 24 instances divide partitions evenly. The team picks 48 with replication factor 3: 144 replicas across six brokers, or 24 per broker, which is negligible. The mistake to avoid is the opposite reflex of 1,000 partitions just in case. Each one adds files, index memory, controller metadata and leader elections on failover, and spreading the same throughput thinly over many partitions makes producer batches smaller and less efficient.
def partitions_needed(target_mb_s, per_partition_produce_mb_s, per_consumer_mb_s, growth=2.0):
base = max(target_mb_s / per_partition_produce_mb_s, target_mb_s / per_consumer_mb_s)
need = base * growth
for n in (6, 12, 24, 36, 48, 60, 72, 96, 120): # divisor-friendly counts
if n >= need:
return n
return int(need)
partitions_needed(60, 25, 4) # -> 36: 15 x 2 = 30, next divisor-friendly countThe function returns 36. The team rounds up to 48 because a second consumer group with heavier per-record work is planned. Write down the rationale in the topic's configuration review, because the number will be questioned later.
Changing the count: add only, and keys move
You can increase a topic's partition count with kafka-topics.sh --alter --partitions N; you cannot decrease it. Increasing it does not move existing data. It changes the result of the key mapping, which is a hash of the key modulo the partition count. After going from 12 to 16 partitions, most keys hash to a different partition than before, while their history stays in the old one. A consumer keeping per-key state, a stream join, or anything relying on per-key order now sees a key's old and new records in different partitions with no ordering between them.
Safe options are to overprovision at creation, to add partitions only for topics without key semantics, or to migrate: create a new topic with the new count, dual-write or mirror into it, move consumers once the old topic drains, then retire it. Treat a partition-count change on a keyed topic as a schema migration. The same mapping concern appears in database sharding, where modulo placement makes resharding expensive.
Reassignment: moving replicas without hurting traffic
Adding brokers does not move existing partitions; new brokers only receive new topics until you reassign. Reassignment copies a partition's full directory to the target brokers, waits until they join the in-sync set, then removes the old replicas. For a 500 GB partition, that copy competes with production traffic for disk and network, so it must be throttled.
# 1. propose a plan for moving these topics onto brokers 1-6
cat > topics.json <<'EOF'
{"version":1,"topics":[{"topic":"payments"}]}
EOF
bin/kafka-reassign-partitions.sh --bootstrap-server kafka:9092 \
--topics-to-move-json-file topics.json --broker-list "1,2,3,4,5,6" --generate > plan.txt
# plan.txt holds the current and the proposed assignment; save the
# "Proposed partition reassignment configuration" JSON as plan.json
# (keep the current one too, as your rollback plan)
# 2. execute with a replication throttle (bytes/sec), then poll
bin/kafka-reassign-partitions.sh --bootstrap-server kafka:9092 \
--reassignment-json-file plan.json --execute --throttle 50000000
bin/kafka-reassign-partitions.sh --bootstrap-server kafka:9092 \
--reassignment-json-file plan.json --verify # also removes the throttle when doneMove a few partitions at a time, watch under-replicated partitions and produce latency, and always run --verify to completion. A throttle left in place silently slows later recovery. Tools such as Cruise Control automate this planning; the mechanics underneath are the same.
Failure modes
| Symptom | Partition-level cause | Response |
|---|---|---|
| Disk full despite short retention | Active segments never roll on quiet partitions; deletion needs closed segments | Lower segment.ms on those topics; check for stuck compaction |
| Consumer lag on a few partitions only | Skewed keys concentrate load on one leader | Fix key choice or add a salt; more partitions do not help a single hot key |
| Slow broker restart | Many partitions to recover and index; leaders to elect | Fewer, larger partitions; controlled shutdown; keep partition count per broker in check |
| Writes rejected with NotEnoughReplicas | In-sync set fell below min.insync.replicas | Find the slow or dead follower; never lower the minimum to make the error go away |
| Per-key order broken after scaling | Partitions added to a keyed topic | Migrate to a new topic instead of altering |
| Cross-zone network bill | All consumers fetch from leaders | Enable rack-aware follower fetching |
What to do next
- For each important topic, write down its measured per-partition producer and per-consumer throughput, and recompute the partition count with headroom.
- Set
broker.rackon every broker and confirm withkafka-topics.sh --describethat no partition has two replicas in one zone. - Use replication factor 3,
min.insync.replicas=2andacks=allfor data you cannot lose; keep unclean leader election disabled. - Find low-traffic topics whose
segment.msexceeds their retention and lower it. - Ban in-place partition increases on keyed topics in your runbook; document the migrate-to-new-topic procedure.
- Rehearse a throttled reassignment in staging and record how long a typical partition takes to move.
- If you run in several zones, enable follower fetching and compare cross-zone traffic before and after.
- Read the Kafka streaming architecture overview to place these choices in the wider pipeline.