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.

Advertisement

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.

One partition: a replicated directory of segment filesLeader replicabroker 1, rack aFollower replicabroker 2, rack bFollower replicabroker 3, rack cfetchorders-7/ on broker 1.logbase 0.logbase 41210active .logbase 83355.indexoffset to byte.timeindextime to offsetepoch checkpointepoch to startOffsets in the active segmentlog start ... HW = 90120 ... LEO = 90155consumers read below HW (committed)records above HW wait for in-sync followersRetentiondeletes whole closed segmentsCompactionrewrites closed segments per keyReassignmentcopies the directory to new brokersProducers append only to the leader's active segment; followers fetch and append the same bytes.The high watermark advances as in-sync followers catch up; only records below it are visible to consumers.Offsets and broker layout are illustrative
A partition is a directory on each replica's broker. Closed segments are immutable; the active segment takes appends. Index files map offsets and timestamps to byte positions. The high watermark separates committed records, visible to consumers, from records still replicating.

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)
Advertisement

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 count

The 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 done

Move 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

SymptomPartition-level causeResponse
Disk full despite short retentionActive segments never roll on quiet partitions; deletion needs closed segmentsLower segment.ms on those topics; check for stuck compaction
Consumer lag on a few partitions onlySkewed keys concentrate load on one leaderFix key choice or add a salt; more partitions do not help a single hot key
Slow broker restartMany partitions to recover and index; leaders to electFewer, larger partitions; controlled shutdown; keep partition count per broker in check
Writes rejected with NotEnoughReplicasIn-sync set fell below min.insync.replicasFind the slow or dead follower; never lower the minimum to make the error go away
Per-key order broken after scalingPartitions added to a keyed topicMigrate to a new topic instead of altering
Cross-zone network billAll consumers fetch from leadersEnable rack-aware follower fetching

What to do next

  1. For each important topic, write down its measured per-partition producer and per-consumer throughput, and recompute the partition count with headroom.
  2. Set broker.rack on every broker and confirm with kafka-topics.sh --describe that no partition has two replicas in one zone.
  3. Use replication factor 3, min.insync.replicas=2 and acks=all for data you cannot lose; keep unclean leader election disabled.
  4. Find low-traffic topics whose segment.ms exceeds their retention and lower it.
  5. Ban in-place partition increases on keyed topics in your runbook; document the migrate-to-new-topic procedure.
  6. Rehearse a throttled reassignment in staging and record how long a typical partition takes to move.
  7. If you run in several zones, enable follower fetching and compare cross-zone traffic before and after.
  8. Read the Kafka streaming architecture overview to place these choices in the wider pipeline.
Key takeaway: A Kafka partition is a replicated directory of immutable segments plus one active segment, indexed sparsely by offset and time, with a high watermark separating committed from replicating records. Retention and compaction work on whole closed segments, placement is per replica and rack-aware only if you configure it, and the count can only grow, remapping keys when it does. Size partitions from measured consumer throughput with headroom, place replicas across zones, throttle every reassignment, and treat any partition-count change on a keyed topic as a migration.