Skip to main content
Create your own

Cassandra's Eventual Consistency: N, W, and R

Hello! Welcome back to our module on Cassandra.

In the last lesson, we established how to configure data replication across multiple datacenters using NetworkTopologyStrategy and defined the Replication Factor (N)—the total number of data copies. This setup ensures your data is resilient to failures.

Today, we address the next logical question: with N copies of our data available, how do we interact with them during read and write operations? This brings us to one of Cassandra's most powerful and defining features: tunable consistency.

This lesson covers the learning outcome: Explain Cassandra's eventual consistency model, including the roles of replication factor (N) and tunable consistency levels (W, R).

We will explore:

  • The core concepts of tunable consistency and why Cassandra is considered an "eventually consistent" system.
  • The different consistency levels available for reads (R) and writes (W).
  • The formula (R + W > N) for achieving strong consistency and its practical implications.
  • The background mechanisms Cassandra uses to ensure data converges across replicas over time.

Understanding these trade-offs is fundamental to designing high-performance, resilient systems. Your work in low-latency trading and high-load payments involves constant balancing of consistency, availability, and performance, and this lesson provides the vocabulary and mental model for managing those trade-offs in Cassandra.

1. The Foundation: Eventual and Tunable Consistency

Cassandra is designed for high availability and partition tolerance, placing it in the "AP" category of the CAP theorem. By default, this means it prioritizes keeping the system online over ensuring that every replica has the identical, latest version of data at all times. This is known as eventual consistency.

However, Cassandra's key innovation is that it doesn't force a single consistency model. Instead, it offers tunable consistency, allowing you to decide the consistency guarantee you need on a per-operation basis.

This is controlled by three variables:

  • N (Replication Factor): The total number of replicas for a given piece of data. You define this per keyspace.
  • W (Write Consistency Level): The number of replica nodes that must acknowledge a write operation for it to be considered successful.
  • R (Read Consistency Level): The number of replica nodes that must respond to a read operation for it to be considered successful.

Let's explore these concepts in more detail.

Replication and consistency

To start, let's get a formal understanding of replication and consistency from the DataStax documentation. This resource provides a clear, structured overview of the key concepts.

Please read the sections "Replication" and "Consistency level", up to but not including the table of write consistency levels. Focus on how the replication factor (N) and the consistency level work together to determine the success of an operation.

As the reading explains, for any write, the coordinator node attempts to send the data to all N replicas. The write consistency level W simply determines how many acknowledgements the coordinator must wait for before it reports success to the client. The remaining writes happen in the background.

This is visualized well in the following image, which compares a write with a higher consistency level (LOCAL_QUORUM) to one with the lowest (ONE).

This diagram illustrates two write scenarios with a Replication Factor of 3. In the top scenario (CL=LOCAL_QUORUM), the coordinator waits for acknowledgements from two replicas before confirming the write. In the bottom scenario (CL=ONE), it waits for only one, significantly reducing latency but offering weaker guarantees.

2. A Spectrum of Consistency Levels

Cassandra offers a wide range of consistency levels for both reads and writes, allowing you to choose the right point on the spectrum from maximum availability to maximum consistency.

Replication and consistency

The DataStax documentation provides comprehensive tables detailing each consistency level. Let's study them to understand the options available.

Please review the tables under "Write consistency" and "Read consistency". You don't need to memorize every level, but focus on understanding the purpose of these key levels: ONE, QUORUM, LOCAL_QUORUM, and ALL.

Let's categorize these levels based on their typical use cases:

  • Maximum Availability (ONE, ANY):

    • ONE: Requires an acknowledgement from only one replica. This offers the lowest latency and highest availability, as the operation can succeed even if most replicas are down. The risk is that a subsequent read might hit a different replica that hasn't yet received the write, resulting in a stale read.
    • ANY (writes only): The write is successful if any node accepts it, even if it's just a hint stored on the coordinator because all replicas are down. This provides the highest availability but the weakest durability guarantee; if the coordinator fails before the hint is delivered, the write is lost.
  • Maximum Consistency (ALL):

    • Requires an acknowledgement from all N replicas. This guarantees that a read will always see the latest write, but it has the highest latency and lowest availability. If even one replica node is down or slow, the operation will fail.
  • The Balanced Approach (QUORUM levels):

    • A quorum is a majority of replicas: floor(N / 2) + 1. For N=3, a quorum is 2. For N=5, a quorum is 3.
    • QUORUM: Requires a quorum of replicas across all datacenters to respond. This is a common choice for strong consistency guarantees.
    • LOCAL_QUORUM: Requires a quorum of replicas within the local datacenter (the datacenter of the coordinator node). This is extremely useful in multi-DC setups, as it provides strong consistency within a single DC without incurring cross-continent latency for every operation.
    • EACH_QUORUM: Requires a quorum in each datacenter to respond. This is used for strong global consistency but fails if any datacenter cannot achieve a quorum.

3. Achieving Strong Consistency: The R + W > N Rule

You can achieve strong consistency—meaning a read is guaranteed to see the result of the last successful write—by tuning your R and W values. The formula is simple but powerful:

Why does this work? It's a pigeonhole principle argument. If the number of nodes you read from (R) plus the number of nodes you wrote to (W) is greater than the total number of replicas (N), your read and write sets must have at least one node in common. This overlapping node guarantees that your read will see the latest write.

Let's see this in action with N=3:

  • Strong Consistency: W=QUORUM(2) and R=QUORUM(2).

    • 2 + 2 = 4, which is greater than 3. This configuration guarantees strong consistency and is a very common production setup. It can tolerate the failure of one replica.
  • Eventual Consistency: W=ONE(1) and R=ONE(1).

    • 1 + 1 = 2, which is not greater than 3. There is no guaranteed overlap. A read might hit one of the two replicas that did not participate in the initial write acknowledgement, returning stale data.
  • Tuning for Workloads:

    • Write-heavy: W=ONE(1), R=ALL(3). 1 + 3 > 3. Writes are fast, but reads are slow and expensive, yet still strongly consistent.
    • Read-heavy: W=ALL(3), R=ONE(1). 3 + 1 > 3. Writes are slow and have low availability, but reads are very fast and strongly consistent.

Replication and consistency

The DataStax article concludes by explaining this formula and providing examples.

Please read the final section, "Immediate consistency", which formalizes the R + W > N rule and shows a table of common combinations.

4. How "Eventual" is Eventual Consistency?

Even when you use low consistency levels, Cassandra has several background mechanisms designed to ensure all replicas eventually converge to the same state.

Eventual Consistency in Apache Cassandra

This article from Medium provides an excellent, practical discussion of Cassandra's internal consistency mechanisms and real-world behavior. It answers the crucial question: how long does 'eventual' actually take?

Please read the following sections: "Introduction" and "Tunable Consistency Levels in Cassandra" for a good recap. "Internal Mechanisms for Eventual Consistency": Pay close attention to the descriptions of Hinted Handoff, Read Repair, and Anti-Entropy Repair. "How Fast Does Cassandra Converge in Practice?" for a practical perspective on convergence times. "Best Practices" and "Conclusion" for a summary of how to apply these concepts.

Let's summarize the key background processes from the article:

  1. Hinted Handoff: If a replica node is down during a write, the coordinator stores a "hint" locally. When the downed node comes back online, the coordinator delivers the hint, allowing the node to catch up on the missed write. By default, hints are stored for 3 hours, covering most transient failures.

  2. Read Repair: When you perform a read with CL > ONE, the coordinator requests data from multiple replicas. If it detects that some replicas are stale (by comparing timestamps), it returns the most recent version to the client and asynchronously triggers a write to update the stale replicas. This is a powerful self-healing mechanism.

  3. Anti-Entropy Repair: This is a maintenance process (run via nodetool repair) that compares all data between replicas and synchronizes any differences. It uses Merkle trees to efficiently find inconsistencies. Running repairs regularly (e.g., weekly, within the gc_grace_seconds window) is a critical operational task to prevent permanent data divergence and fix issues that hints and read repair might miss, especially related to deletes (tombstones).

In a healthy cluster, "eventual" is typically on the order of milliseconds. The window of inconsistency primarily exists during network partitions or node failures.

Conclusion

In this lesson, we demystified Cassandra's eventual consistency model, revealing it to be a flexible and powerful system of tunable guarantees.

Key Takeaways:

  • Cassandra's consistency is tunable on a per-operation basis using Write (W) and Read (R) consistency levels.
  • The formula R + W > N is the key to achieving strong consistency, where N is the Replication Factor.
  • Quorum-based levels (QUORUM, LOCAL_QUORUM) provide a robust and popular balance between consistency, performance, and availability.
  • Even with low consistency settings, Cassandra works to converge data using background processes like hinted handoff, read repair, and anti-entropy repair. Regular nodetool repair is a crucial operational responsibility.

Preview of the Next Lesson:

We have now covered how data is partitioned across nodes, replicated for fault tolerance, and accessed with tunable consistency. The next step is to look inside a single Cassandra node. How does it manage writes and reads on disk to deliver such high performance?

In our next lesson, we will explore Cassandra's storage architecture, focusing on the Log-Structured Merge-Tree (LSM-Tree), Memtables, and SSTables. This will explain the mechanics behind Cassandra's famously fast write operations.

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

Sign up