Learn Labs
6. Kafka Internals

6.4 Replication

Scale note: "each broker typically stores hundreds or even thousands of replicas belonging to different topics and partitions."

ISR membership

leaderin sync
follower 1in sync
follower 2in sync
3/3
ISR
2
required
Safe

All 3 replicas are in sync. Committed messages exist on every one of them.

Since 0.9 the ISR is defined by time, not message count: a replica is in sync if it has requested the leader's latest message within replica.lag.time.max.ms. That single change removed the tuning guesswork around lag in messages.

"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 replicaFollower replica
Count per partitionExactly oneAll the rest
Produce requestsAll produce requests go through the leader — to guarantee consistency—
Consume requestsYes"Unless configured otherwise, followers don't serve client requests"
Main jobServe 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 leader
  • RackAwareReplicaSelector — match the broker's rack.id to client.rack
  • your own ReplicaSelector implementation

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."

HW piggybacked on the replication responseleader[ …committed… | HW ]leader-consumersavailable NOWfollower[ …committed… | HW (older) ]follower-consumersavailable LATER

Follower reads are cheaper (no cross-AZ egress) but slightly staler. You are trading latency for cost, deliberately.

Figure 6.4.2

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,2

Broker 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.


On this page