Skip to main content
Create your own

Kafka ISR Architecture and Consistency

Hello! Welcome back to your course on high-load distributed systems.

In our last lesson, we configured Kafka topics using replication-factor to create multiple copies of our data for fault tolerance. However, simply creating copies isn't enough. We need a robust mechanism to manage these copies, ensure they are consistent, and use them to recover from failures without losing data.

Today, we will explore the architectural heart of Kafka's durability and consistency model: the In-Sync Replica (ISR) set. Understanding the ISR is crucial for reasoning about Kafka's behavior during failures and for making informed decisions about data safety guarantees.

Lesson Goal: By the end of this session, you will be able to explain the architecture of Kafka's In-Sync Replicas (ISR) and its critical role in providing consistency.

1. From Replicas to In-Sync Replicas

As we've established, a topic partition has one leader replica and one or more follower replicas. All produce and consume requests are served by the leader. The followers have a simple job: copy everything the leader does.

This raises a critical question: if the leader fails, which follower is eligible to become the new leader? The only safe choice is a follower that is perfectly caught up with the leader. A follower that is lagging behind might not have the latest messages, and promoting it would lead to data loss.

This is where the In-Sync Replica (ISR) set comes in. It is the set of replicas, including the leader itself, that are considered fully "caught up."

This diagram shows the basic replication flow. A producer writes to the leader (Broker 0). The leader replicates the data to followers (Broker 1, Broker 2). A message is only considered "committed" and visible to consumers after it has been successfully copied to the replicas in the ISR.

To understand the formal definition and mechanics of the ISR, let's turn to the Confluent documentation.

Kafka Replication and Committed Messages

The Confluent documentation 'Kafka Replication and Committed Messages' provides a precise definition of the ISR and how Kafka determines if a replica is in sync.

Please read the sections 'Replication overview' and 'In-sync replicas and producer acks'. Focus on the two conditions for a node to be considered 'alive' and how the ISR set is defined. Pay close attention to the definition of a 'committed message'.

As you read, Kafka's definition of "in-sync" is not just about being connected. A follower must be actively fetching data from the leader and not fall too far behind. This lag is controlled by the broker-level configuration replica.lag.time.max.ms. If a follower fails to send a fetch request or falls behind the leader's log for more than this duration, the leader removes it from the ISR.

2. The High Watermark: The Commitment Boundary

The reading introduced the concept of a "committed message." This is the cornerstone of Kafka's consistency for consumers. A consumer will never be shown a message that is not yet committed. This prevents a scenario where a consumer reads a message that is later lost due to a leader failure before it could be replicated.

Kafka tracks the boundary between committed and uncommitted messages using a special offset called the High Watermark (HWM).

The process works as follows:

  1. A producer sends a message to the partition leader.
  2. The leader appends the message to its own log.
  3. Followers in the ISR fetch the new message from the leader.
  4. Once all replicas in the ISR have written the message to their logs, the leader advances the High Watermark offset.
  5. Only messages at or below the High Watermark are considered committed and are made available to consumers.

Let's explore this mechanism in more detail.

Kafka Data Replication Protocol: A Complete Guide

The Confluent course material 'Kafka Data Replication Protocol' explains the mechanics of the High Watermark and how it's advanced.

Please read the section 'Committing Partition Offsets' and 'Advancing the Follower High Watermark'. Focus on how the leader tracks the progress of its followers and decides when to advance the HWM.

This mechanism ensures that a message is only visible to consumers after it has been durably persisted across the required set of machines, as defined by the current ISR.

This diagram illustrates the write-commit flow when a producer requires full acknowledgment (`acks=all`). The write goes to the leader, which replicates it to followers in the ISR. Only after receiving acknowledgments from all ISR followers does the leader 'commit' the data (i.e., advance the High Watermark) and acknowledge the write to the producer.

3. ISR Dynamics and Fault Tolerance

The ISR is not static; it's a dynamic set that changes based on the health and performance of the cluster. This dynamism is what gives Kafka its resilience.

Case 1: A Follower Fails or Lags

If a follower broker crashes or becomes too slow to keep up, the leader will remove it from the ISR after the replica.lag.time.max.ms timeout. The ISR shrinks, but the system continues to operate. The leader can still commit new messages with the remaining replicas in the ISR. This ensures that a single slow or failed replica does not halt the entire system.

Kafka Data Replication Protocol: A Complete Guide

Let's look at how Kafka handles a lagging follower.

Please read the short section 'Handling Failed or Slow Followers'. This explains the process of removing a follower from the ISR.

Case 2: The Leader Fails

This is the critical failure scenario. When the leader broker becomes unavailable, the Kafka controller must elect a new leader.

  • The new leader is chosen from the ISR. Because every replica in the ISR is guaranteed to have all committed messages (up to the HWM), any one of them can be promoted to leader without losing data.
  • This provides Kafka's core data loss guarantee: as long as there is at least one replica remaining in the ISR, committed data will not be lost.

Case 3: A Replica Recovers

When a broker that was previously down comes back online, its partition replicas are out of sync. It cannot immediately rejoin the ISR. It must first go through a reconciliation process to catch up with the new leader. This involves a crucial step: log truncation. The recovering replica may have old, divergent data if it was a leader that failed. It must truncate its log to match the new leader's log and then fetch all the new data before it can be added back into the ISR.

Given your background, you might find this low-level reconciliation process particularly interesting.

Kafka Data Replication Protocol: A Complete Guide

The 'Kafka Data Replication Protocol' guide details the log reconciliation process that a recovering replica must undergo.

Please read the section 'Partition Replica Reconciliation'. This explains how an out-of-sync follower uses leader epoch information to truncate its log and get back in sync.

4. The Ultimate Trade-Off: Consistency vs. Availability

We've established that Kafka guarantees no loss of committed data as long as at least one ISR member survives. But what happens in the catastrophic event that all replicas for a partition fail, or the only surviving replica is one that was not in the ISR?

This forces a fundamental choice in distributed systems: do you prioritize consistency (guaranteeing you don't lose data) or availability (allowing the system to accept writes)?

Kafka makes this choice configurable.

Kafka Replication and Committed Messages

The Confluent documentation on replication discusses this critical trade-off and the configuration that controls it.

Please read the section 'Unclean leader election and partition loss'. This explains the two possible behaviors when all ISR replicas are lost and the role of the unclean.leader.election.enable setting.

  • unclean.leader.election.enable = false (Default): Prioritize Consistency. Kafka will wait for a replica from the original ISR to come back online to elect it as the new leader. The partition remains unavailable for writes until then. If the data for all ISR members is permanently lost, the partition is lost forever.
  • unclean.leader.election.enable = true: Prioritize Availability. Kafka will elect the first replica to come back online as the new leader, even if it was not in the ISR and has therefore lost data. This makes the partition available again but at the cost of sacrificing consistency.

This is a critical operational decision that directly reflects business priorities. For a financial system like a payment platform, you would almost certainly favor consistency. For a system processing less critical data like metrics or logs, availability might be more important.

A Note on Quorums

Your background in distributed systems theory might bring to mind consensus algorithms like Paxos or Raft, which often rely on a majority quorum. Kafka's ISR model is a different approach. Instead of a static majority, it uses a dynamic quorum (the ISR set) that can shrink or expand. This allows Kafka to remain available for writes even if a majority of replicas fail, as long as the minimum number of in-sync replicas is met. This flexibility is advantageous in a system designed to manage tens of thousands of partitions.

Conclusion

In this lesson, we dissected the In-Sync Replica (ISR) mechanism, the engine that powers Kafka's replication, consistency, and durability guarantees.

Key Takeaways:

  • The ISR is a dynamic set of replicas (leader and followers) that are fully caught up with the partition leader.
  • A message is committed only after it has been written to all replicas in the ISR. This point is tracked by the High Watermark (HWM).
  • Consumers can only read committed messages, preventing them from seeing data that might be lost in a failure.
  • Fault-tolerant leader election relies on choosing a new leader from the ISR, guaranteeing no loss of committed data.
  • Kafka allows you to make a direct trade-off between consistency and availability via the unclean.leader.election.enable configuration for catastrophic failure scenarios.

Next Up

Now that you have a firm grasp of the server-side guarantees provided by the ISR, our next lesson will focus on the client side. We will explore the producer's acks configuration (acks=0, acks=1, acks=all) and analyze how each setting interacts with the ISR to provide different trade-offs between write latency and durability.

Can't find a good explanation? Sign up and we'll make it for you

Sign up