# How Data Consistency Works in Distributed Message Queues with Replication and ISR

> Learn how distributed message queues ensure data consistency using replication and ISR. Discover how In Sync Replicas guarantee committed records survive failures.

- Repository: [Gaurav Kumar/system-design-notes](https://github.com/liquidslr/system-design-notes)
- Tags: deep-dive
- Published: 2026-09-11

---

**Data consistency in distributed message queues is enforced through a leader-follower replication model where writes are acknowledged only after all In‑Sync Replicas (ISR) persist the message, ensuring committed records survive single-node failures.**

The `liquidslr/system-design-notes` repository documents a reference architecture for distributed message queues that mirrors production systems like Apache Kafka. This design achieves **data consistency in distributed message queues with replication and ISR** by combining quorum-based replication, strict commit boundaries, and coordinator-managed metadata.

## Replication Architecture and the Leader-Follower Model

Distributed message queues partition topics across multiple brokers to achieve horizontal scalability. Each partition is replicated across a configurable number of nodes to ensure fault tolerance.

### Partition Replication Strategy

According to the source documentation in `19. Distributed Message Queue/README.md` (lines 78‑90), the system implements the following replication protocol:

- Each **partition** is stored on multiple broker nodes.
- One replica is elected **leader**; remaining replicas serve as **followers**.
- Producers write exclusively to the leader.
- Followers continuously pull new log entries from the leader to maintain synchronization.

The **replica distribution plan**—which broker holds which replica—is stored in **Zookeeper** and updated dynamically whenever brokers join or leave the cluster. This ensures the routing layer always directs writes to the current leader.

### Metadata Management in Zookeeper

Zookeeper provides the strong consistency guarantees required for coordination. As documented in lines 58‑67, it stores:

- **Metadata**: Topic configurations and replica assignments.
- **State**: Consumer offsets and the current ISR list.
- **Quorum consistency**: Ensures all brokers share an identical view of the cluster topology through Zookeeper’s internal consensus protocol.

## In‑Sync Replicas (ISR) and Commit Semantics

The **In‑Sync Replicas (ISR)** set defines the subset of followers that are sufficiently caught up with the leader to participate in the commit process.

### Defining the ISR Set

Per lines 94‑101 of the README, a follower remains in the ISR if it satisfies the `replica.lag.max.messages` threshold. This configurable lag bound determines whether a replica is considered "caught up." The leader maintains the ISR list in memory and persists it through Zookeeper.

### The Commit Boundary

The leader advances the **committed offset** only after a message has been replicated to all ISR members. This creates a strict durability boundary:

1. The leader writes the message to its local log (WAL segment).
2. The message is propagated to all followers currently in the ISR.
3. The leader waits for acknowledgments from every ISR replica.
4. Only then is the offset marked as committed and made visible to consumers.

This mechanism ensures that committed messages exist on a quorum of nodes, specifically all in-sync replicas, rather than just the leader.

## Producer Acknowledgment Modes and Consistency Guarantees

Producers control durability guarantees through the `requiredAcks` parameter. The system supports three acknowledgment modes (documented in lines 16‑30):

| Ack mode | Producer wait condition | Consistency impact |
|----------|------------------------|-------------------|
| `ACK=all` | All ISR replicas persist the record | **Strong durability**—committed messages survive any single replica failure |
| `ACK=1` | Only the leader persists the record | Higher throughput risk: leader crash before followers catch up causes loss |
| `ACK=0` | No acknowledgment | Maximum throughput, zero durability guarantee |

When `ACK=all` is configured, the leader invokes `waitUntilAllISRAck()` before advancing the committed offset, blocking the producer until the ISR quorum acknowledges persistence.

## Failure Recovery and Leader Election

The system maintains availability through automated failover procedures coordinated by Zookeeper.

### Handling Broker Failures

When a broker fails, Zookeeper detects the missed heartbeat and triggers leader election. As described in lines 61‑71:

1. A surviving replica from the current ISR is promoted to leader.
2. Zookeeper updates the metadata to reflect the new leader assignment.
3. The routing layer redirects producer traffic to the new leader.

Only replicas that were in the ISR at the moment of failure are eligible for leadership, ensuring that the new leader contains all committed messages.

### Re‑synchronization Process

After a failed broker recovers, it rejoins as a follower and begins a **re‑synchronization** phase:

- The follower truncates its local log to the last committed offset known to the new leader.
- It then pulls all missing messages from the leader.
- Once the lag falls below `replica.lag.max.messages`, the follower is re‑added to the ISR set.

This process maintains the invariant that all ISR members contain identical log segments up to the committed offset.

## Implementation Details from the Source

The following pseudo‑code illustrates the consistency protocol as implemented in the reference architecture.

### Producer Write Path

```pseudo
function produce(topic, partitionKey, payload):
    leader = routingLayer.getLeader(topic, partitionKey)
    msg = createMessage(payload, partitionKey)
    ack = leader.append(msg, requiredAcks = "all")
    if ack.success:
        return SUCCESS
    else:
        retryOrFail()

```

### Leader Append and Commit Logic

```pseudo
function append(msg, requiredAcks):
    writeToLocalLog(msg)
    for follower in ISR:
        sendReplication(msg, follower)
    
    if requiredAcks == "all":
        waitUntilAllISRAck(msg)
        commitOffset(msg.offset)
        return ACK_SUCCESS
    else if requiredAcks == "1":
        commitOffset(msg.offset)
        return ACK_SUCCESS

```

### Follower Replication Loop

```pseudo
function replicationLoop():
    while true:
        msg = receiveFromLeader()
        writeToLocalLog(msg)
        sendAckToLeader(msg.offset)

```

## Summary

- **Leader-follower replication** ensures writes are ordered and replicated across multiple brokers, as defined in `19. Distributed Message Queue/README.md` (lines 78‑90).
- **In‑Sync Replicas (ISR)** establish the minimum replica set that must acknowledge a write before it is considered committed (lines 94‑101).
- **Zookeeper coordination** provides strongly consistent metadata storage for replica assignments and leader election (lines 58‑67).
- **Configurable acknowledgment modes** allow producers to trade latency for durability, with `ACK=all` providing the strongest consistency guarantees (lines 16‑30).
- **Automatic failover** promotes an ISR member to leader during broker failures, ensuring committed messages remain available (lines 61‑71).

## Frequently Asked Questions

### What happens if a follower falls behind the leader?

A follower that exceeds the `replica.lag.max.messages` threshold is removed from the ISR set. The leader stops waiting for acknowledgments from this replica before committing new messages. Once the follower catches up and acknowledges the backlog, it is re‑added to the ISR and resumes participation in the commit quorum.

### How does ACK=all differ from ACK=1 in terms of data safety?

With `ACK=all`, the producer receives confirmation only after all ISR replicas persist the message, guaranteeing that a committed write survives any single broker failure. With `ACK=1`, the producer receives confirmation after only the leader writes the message, creating a window of vulnerability where a leader crash before followers replicate the data results in permanent message loss.

### What role does Zookeeper play during a leader failure?

Zookeeper maintains the authoritative registry of broker memberships and ISR lists. When a leader fails, Zookeeper coordinates the election of a new leader from the remaining ISR members and updates the cluster metadata atomically. This ensures that all brokers and clients converge on a consistent view of the new topology without split-brain scenarios.

### Can a message be lost after being committed to the ISR?

No committed message can be lost as long as at least one ISR member survives. Because the commit offset advances only after all ISR replicas acknowledge the message, a minimum of one replica (including the leader) holds the data at commit time. Subsequent failures cannot remove the message unless all ISR members fail simultaneously before the next checkpoint.