6.4 Replication
Scale note: "each broker typically stores hundreds or even thousands of replicas belonging to different topics and partitions."
ISR membership
- 3/3
- ISR
- 2
- required
All 3 replicas are in sync. Committed messages exist on every one of them.
"Replication is at the heart of Kafka's architecture. Indeed, Kafka is often described as 'a distributed, partitioned, replicated commit log service.' Replication is critical because it is the way Kafka guarantees availability and durability when individual nodes inevitably fail."
Scale note: "each broker typically stores hundreds or even thousands of replicas belonging to different topics and partitions."
4.1 Leader and follower replicas
| Leader replica | Follower replica | |
|---|---|---|
| Count per partition | Exactly one | All the rest |
| Produce requests | All produce requests go through the leader — to guarantee consistency | — |
| Consume requests | Yes | "Unless configured otherwise, followers don't serve client requests" |
| Main job | Serve clients; track follower progress | "replicate messages from the leader and stay up-to-date with the most recent messages the leader has" |
| On leader crash | — | "one of the follower replicas will be promoted to become the new leader" |
4.2 Read from follower (KIP-392) — and its hidden latency cost
Goal: "decrease network traffic costs by allowing clients to consume from the nearest in-sync replica rather than from the lead replica."
Consumer config: client.rack = the location of the client.
Broker config: replica.selector.class —
LeaderSelector— default, always the leaderRackAwareReplicaSelector— match the broker'srack.idtoclient.rack- your own
ReplicaSelectorimplementation
How correctness is preserved — and why it costs latency:
"The replication protocol was extended to guarantee that only committed messages will be available when consuming from a follower replica... To provide this guarantee, all replicas need to know which messages were committed by the leader. To achieve this, the leader includes the current HIGH-WATER MARK (latest committed offset) in the data that it sends to the follower."
⚠️ "The propagation of the high-water mark introduces a small delay, which means that data is available for consuming from the leader EARLIER than it is available on the follower. It is important to remember this additional delay, since it is tempting to attempt to decrease consumer latency by consuming from the leader replica."
Follower reads are cheaper (no cross-AZ egress) but slightly staler. You are trading latency for cost, deliberately.
4.3 The ISR mechanism — how "in sync" is actually determined
The elegant part: followers use exactly the same Fetch requests consumers use.
"To stay in sync with the leader, the replicas send the leader Fetch requests — the exact same type of requests that consumers send in order to consume messages... Those Fetch requests contain the offset of the message that the replica wants to receive next, and will always be in order. This means that the leader can know that a replica got all messages up to the last message that the replica fetched, and none of the messages that came after. By looking at the last offset requested by each replica, the leader can tell how far behind each replica is."
The leader needs No extra protocol to track follower progress. A follower asking for offset N is Proof it has everything < N.
⇒ “replication is just consumption” — one mechanism, two uses.
The out-of-sync rule — two conditions:
A replica is considered out of sync if:
- it hasn't requested a message in more than 10 seconds, or
- it has requested messages but hasn't caught up to the most recent message in more than 10 seconds
Controlled by replica.lag.time.max.ms.
Consequence: “If a replica fails to keep up with the leader, it can no longer become the new leader in the event of failure — after all, it does not contain all the messages.”
"replicas that are consistently asking for the latest messages are called in-sync replicas. Only in-sync replicas are eligible to be elected as partition leaders in case the existing leader fails."
Why followers fall behind (the book's examples): "network congestion slows down replication", or "a broker crashes and all replicas on that broker start falling behind until we start the broker and they can start replicating again."
4.4 Preferred leader
"each partition has a preferred leader — the replica that was the leader when the topic was originally created. It is preferred because when partitions are first created, the leaders are balanced among brokers. As a result, we expect that when the preferred leader is indeed the leader for all partitions in the cluster, load will be evenly balanced between brokers."
Default auto.leader.rebalance.enable=true "will check if the preferred leader replica is not the current leader but is in sync, and will trigger leader election to make the preferred leader the current leader."
FINDING THE PREFERRED LEADER
*"The best way to identify the current preferred leader is by looking at the list of replicas for a partition (see
kafka-topics.sh). THE FIRST REPLICA IN THE LIST IS ALWAYS THE PREFERRED LEADER. This is true no matter who is the current leader and even if the replicas were reassigned to different brokers using the replica reassignment tool.In fact, if you manually reassign replicas, it is important to remember that the replica you specify first will be the preferred replica, so make sure you spread those around different brokers to avoid overloading some brokers with leaders while other brokers are not handling their fair share of the work."*
Topic: orders Partition: 5 Replicas: 3,1,2 Isr: 3,1,2Broker 3 — the first replica in the list — is the preferred leader, regardless of who Leader: currently says.
This is the detail that makes Ch. 5's alterPartitionReassignments examples make sense: [1,0] vs [0,1] decides leadership placement. Get it wrong across many partitions and you manufacture a hot broker.
6.3 KRaft — the Raft-based controller
That is the broker-level analogue of controller zombie fencing: a lagging broker can currently accept writes it has no right to accept, because it doesn't yet know it lost leaders…
6.5 Request processing
This is the map for every broker metric and thread-pool config you'll ever tune.