0%
Data-Intensive Applications
Foundations of Data Systems
Encoding and Evolution
Batch Processing
Stream Processing
Data Quality and Governance
Operational Patterns
Consistency and Consensus
When you store data on a single machine, reads always return the latest write. The moment you replicate data across multiple nodes, this guarantee evaporates. A client writes to node A, but a reader on node B sees stale data because the replication has not arrived yet. The question is not whether your data will be inconsistent across replicas. It will be. The question is what promises you make to clients about the inconsistency they might observe.
Consistency guarantees are those promises. They define what a client can expect when reading data that might be in the process of being replicated. Stronger guarantees make applications easier to reason about but cost performance and availability. Weaker guarantees allow higher throughput and lower latency but push complexity into the application.
Understanding these guarantees is not academic. Every time you choose a database, configure a replication mode, or decide how to handle a read after a write, you are making a consistency decision. The wrong choice leads to bugs that are nearly impossible to reproduce in testing because they only appear under specific timing conditions in production: a replica that is 50ms behind, a network packet that arrives out of order, a garbage collection pause that delays a heartbeat.
The Spectrum of Consistency
At the weakest end sits eventual consistency: if no new writes arrive, all replicas will eventually converge to the same value. This says nothing about how long "eventually" takes. It could be milliseconds or hours. A client might read a value, write a new one, then read the old value again. The only guarantee is that silence eventually produces agreement.
One step stronger is read-your-writes (also called read-after-write) consistency: a client that writes a value is guaranteed to see that value on subsequent reads, even if the read hits a different replica. Other clients might still see stale data, but you always see your own writes. This is the minimum viable guarantee for most user-facing applications. Without it, a user saves their profile, refreshes the page, and sees the old profile, a confusing experience that erodes trust in the system.
There are several practical ways to implement read-your-writes. The simplest is sticky sessions: route a user's reads to the same replica that accepted their most recent write. This works but defeats the purpose of load balancing. A better approach is to attach a logical timestamp to the write response. The client sends this timestamp with subsequent reads, and the read path waits until the target replica has advanced past that timestamp before serving the response. The read might be slightly delayed, but it never returns stale data.
Monotonic reads guarantee that once a client sees a value at time T, it will never see an older value on subsequent reads. Time does not go backward. Without this, a client might see a comment, refresh, see the comment gone, then refresh again and see it reappear. This happens when successive reads hit different replicas with different replication lag. The fix is similar to read-your-writes: track the highest timestamp the client has seen and only serve reads from replicas that have reached at least that timestamp.
Consistent prefix reads guarantee that if write A happened before write B, every client that sees B also sees A. This prevents causal violations: seeing a response before its question, or seeing a transfer credit without the corresponding debit. In partitioned databases, this is particularly tricky because different partitions replicate independently. If the question and answer land on different partitions, they may arrive at a reader in the wrong order.
In practice, most applications need read-your-writes plus monotonic reads at minimum. These can be achieved without full linearizability by routing a user's reads to the same replica that handled their writes, or by tracking a logical timestamp and waiting for replicas to catch up before serving reads.
Why Stronger is Not Always Better
Linearizability, the strongest single-object guarantee, makes every operation appear to take effect at a single instant. It is the gold standard for correctness. But it comes with severe costs. Every write must be acknowledged by a majority of replicas before it is visible. Every read must check whether a more recent write exists. Under network partitions, a linearizable system must either refuse operations (sacrificing availability) or risk violating its guarantee.
The CAP theorem formalizes this tradeoff. In a distributed system experiencing a network partition, you must choose between consistency (linearizability) and availability (every non-failing node responds). Since network partitions are inevitable in any system that spans multiple machines, the real question is: when a partition occurs, do you want your system to return errors (CP) or potentially stale data (AP)?
A common misunderstanding of CAP is that you choose two of three properties at design time. This is misleading. When the network is healthy, you have all three: consistency, availability, and partition tolerance. The choice is forced only during a partition, which is a runtime event, not a design-time decision. Furthermore, "availability" in CAP means every non-failing node must respond, which is a very strong definition. Most practical systems operate in a spectrum between CP and AP, choosing different points for different operations.
Systems like Google Spanner push the boundary by using GPS-synchronized clocks (TrueTime) to provide linearizability with high availability. Spanner achieves this by bounding clock uncertainty to a few milliseconds and waiting out that uncertainty window before committing. This is not free: every write pays a latency cost equal to the clock uncertainty, and the system depends on specialized hardware (GPS receivers and atomic clocks in every datacenter).
Mid-level engineers should be able to name the consistency levels and explain why eventual consistency causes user-visible bugs. Senior engineers should be able to choose the right consistency level for a given feature and implement read-your-writes using sticky sessions or logical timestamps. Staff engineers should be able to design systems that mix consistency levels: linearizable for payment balances, eventually consistent for recommendation feeds, and causal for chat messages.
Mixing Consistency Levels in Practice
Most real systems mix consistency levels across different data types. A banking application needs linearizability for account balances but can tolerate eventual consistency for transaction history display. A social network needs causal consistency for comment threads (responses appear after their parent) but can tolerate eventual consistency for like counts.
The architectural pattern is to classify your data by the cost of inconsistency. If stale data causes financial loss (double-spending), use linearizability. If stale data causes user confusion (seeing your old profile after an update), use read-your-writes. If stale data is invisible (a like count that is off by 3 for two seconds), use eventual consistency. Each step down the consistency ladder buys you lower latency, higher throughput, and better availability during partitions. Choosing the right consistency level per data type is one of the most impactful architectural decisions you will make.