The idea: a log that many readers share
Kafka stores messages (records) in topics, and a topic is split into partitions. Each partition is an append-only log: a new record always goes to the end and gets the next offset (0, 1, 2, ...). Records are not deleted when someone reads them; they stay until the retention time or size runs out. A reader is just a position in the log, so many independent applications can read the same data, each at its own pace, and any of them can go back and read it again.
The page runs a topic orders with 3 partitions on 3 brokers, with replication factor 2: every partition has a leader copy and one follower copy on another broker.
| Partition | Leader | Follower |
|---|---|---|
| P0 | B1 | B2 |
| P1 | B2 | B3 |
| P2 | B3 | B1 |
Leaders are spread over the brokers so that each broker does a share of the work. All reads and writes of a partition go to its leader; the follower only keeps a copy in case the leader fails.
Keys, partitions and ordering
The producer's partitioner decides which partition a record goes to. With a key it is murmur2(key) % partitions (the page uses a fixed table instead), so all records with the same key go to the same partition, and Kafka keeps them in the order they were sent. There is no order across partitions: a record in P2 may be read before an older record in P0. Choose the key so that the records whose order matters share it: all events of one user, one account, one order. Without a key, the sticky partitioner fills one partition's batch, then moves to another, which gives big batches and still spreads the load.
The producer: batching and the sender thread
send() does not send anything over the network. It serializes the key and value, picks the partition and appends the record to that partition's current batch in the RecordAccumulator, then returns a future. A separate sender thread takes batches that are ready: full (batch.size, 16 KB by default; 3 records on the page) or old enough (linger.ms; 2 ticks on the page). It groups them by the broker that leads their partition and sends one ProduceRequest per broker, with the batches of several partitions inside. Batching is why Kafka is fast: one request, one disk append and one compression pass cover many records. max.in.flight.requests.per.connection limits the requests that wait for an answer; the page uses 1.
acks: when is a write done?
| acks | The producer hears "done" when | Can a write that was reported done be lost? |
|---|---|---|
| 0 | the request has left; the broker never answers | yes: the producer does not even learn about errors |
| 1 | the leader has appended the batch | yes, if the leader dies before a follower copies it (Demo: leader crash with acks=1) |
| all (-1) | every in-sync replica has it (the high watermark has passed it) | no, as long as at least min.insync.replicas copies exist |
Run Demo: acks=0 vs acks=1 vs acks=all and compare the "sent t, acked t" lines under the sender thread: every step of safety costs a round trip. acks=all is the default since Kafka 3.0.
Replication: followers pull, the high watermark moves
A follower copies its leader by sending Fetch requests, exactly like a consumer does. The fetch offset is also an acknowledgement: "fetch from offset 5" means "I have everything before 5". From these the leader knows each follower's log end offset (LEO), and it computes the high watermark (HW): the smallest LEO among the in-sync replicas (ISR). Records below the HW are on every in-sync copy; they are committed. Only committed records are handed to consumers, so a consumer never sees a record that a leader change could still remove. On the page the HW is the blue bar, and committed cells are coloured by key.
A follower that is down or falls too far behind (replica.lag.time.max.ms, 30 s) is removed from the ISR, so a slow follower cannot hold the HW back forever. It is added again once it has caught up. With acks=all, min.insync.replicas says how small the ISR may get before writes are refused with NOT_ENOUGH_REPLICAS. With replication factor 2 and min.insync.replicas = 2, as on this page, any broker failure stops acks=all writes to the partitions it held; production clusters usually use replication factor 3 with min.insync.replicas = 2, which survives one failure without stopping.
When a leader fails
The controller (a quorum of controller nodes in KRaft mode, ZooKeeper in older versions) watches the brokers. When a leader dies it picks a new leader from the ISR and increases the partition's leader epoch. Producers and consumers get NOT_LEADER_OR_FOLLOWER or a broken connection, refresh their metadata and go to the new leader. When the old leader comes back, it is a follower: it asks the new leader where its own last epoch ended and truncates its log there, dropping any records that were never copied. With acks=1 those records may already have been reported as sent; with acks=all they were not, and the producer sends them again. Run the two leader-crash demos to see both.
Choosing a leader outside the ISR (unclean.leader.election.enable) would keep the partition available but lose committed records; it is off by default.
Consumers pull
A consumer asks the leader for records from its position in each partition. If there is nothing new below the HW, the broker holds the request for up to fetch.max.wait.ms (a long poll) and answers as soon as data arrives, so an idle consumer costs almost nothing and a new record is delivered at once. Pulling lets each consumer go at its own speed and read in big batches; a slow consumer only falls behind, it does not slow the brokers down.
Consumer groups and rebalancing
Consumers with the same group.id share the work: each partition is read by exactly one member of the group, so the number of partitions is the limit on parallelism. With two consumers and three partitions the split is uneven: Demo: 5 records, 2 consumers sends user-1 to user-5, which land on P0 (user-1, user-4), P1 (user-3) and P2 (user-2, user-5); the range assignor gives C1 both P0 and P1, so C1 processes three records and C2 two. How evenly the work spreads depends on the keys and on the number of partitions, not on the number of records. A fourth consumer on a three-partition topic gets nothing (Demo: consumer group scaling). Different groups are independent: each gets every record.
One broker is the group's coordinator. Members send it heartbeats; when a member joins, leaves or misses heartbeats for session.timeout.ms, the coordinator starts a rebalance:
- Every member stops fetching and gives up its partitions (committing its position first), then sends
JoinGroup. - The coordinator starts a new generation and names one member the group leader.
- The leader computes the assignment with the configured assignor (range, round-robin, sticky) and sends it in
SyncGroup; the coordinator gives each member its part. - Each member reads the committed offsets of its new partitions (
OffsetFetch) and starts fetching from there.
This "stop the world" protocol is the eager one the page draws. The cooperative sticky assignor lets members keep the partitions that do not move, and the newer consumer protocol (KIP-848) moves the assignment to the coordinator, but the idea is the same.
Committed offsets and delivery guarantees
A consumer's position lives in its memory. To survive a crash or a rebalance it commits it to the internal topic __consumer_offsets, through the coordinator. The committed offset is the next record to read. When a partition gets a new owner, the new owner starts at the committed offset, not where the old owner stopped, so the time between processing and committing decides what happens on failure:
- Commit after processing (the usual way): a crash between the two means those records are processed again. This is at-least-once, Kafka's default; Demo: consumer crash before commit shows the duplicates. Make processing idempotent (for example an upsert by id) and duplicates do no harm.
- Commit before processing: a crash in between skips those records: at-most-once.
- Exactly-once inside Kafka (read, process, write to another topic) uses an idempotent producer and transactions, which commit the output records and the consumer offsets together. The page does not simulate it; see Kafka Consumer Groups: Rebalancing, Duplicates and Exactly-Once.
The producer has the same problem in the other direction. If a batch is written but the answer is lost, the producer cannot tell and sends it again, and without idempotence the log gets it twice. enable.idempotence=true (the default since Kafka 3.0) gives each producer an id and each batch a sequence number, so the leader drops the repeated copy. The page leaves it out and uses a producer that retries without it.
What the page leaves out
Log segments and retention, log compaction (keeping only the latest record per key), compression of batches, the page cache and zero-copy sendfile() on reads, fetch.min.bytes, heartbeats as separate messages, the controller quorum itself, preferred leader election (moving leadership back after a restart), rack awareness and quotas.