Eventual Consistency

During the partition, shoppers may see the old price of the blue mug. When the connection is restored, the primary sends the replica the changes it missed, and the replica applies them. From then on, a catalog read returns the new price.

This is called eventual consistency: after a write, the replicas can disagree for a while, but they all end up with the same data.

What eventual consistency promises

While the replicas disagree, a read from a replica might return the old data. This can last as long as it takes to restore the lost network connection.

This happens without a partition too. Here is a write with asynchronous replication, which we saw in the previous chapter:

sequenceDiagram
    participant A as Shopend API
    participant P as Primary
    participant R1 as Replica 1
    participant R2 as Replica 2
    A->>P: Write
    P->>P: Record write
    P-->>A: Success
    Note over R1,R2: Reads can still return old data
    P->>R1: Replicate change
    R1->>R1: Apply change
    P->>R2: Replicate change
    R2->>R2: Apply change
    Note over P,R2: All replicas now include the write

The primary reports success before the replicas apply the change. This delay is the replication lag from the previous chapter.

Delivering missed changes and resolving conflicts

To send the replica the changes it missed, the database has to keep those changes until the connection is restored.

In Shopend, only the primary accepts writes, so the replica just applies the primary’s changes. If two replicas accept writes, their changes can conflict. Recall the cart example from earlier in this chapter: failover created a second primary, and each primary accepted one of the shopper’s two edits. Each primary can send its edit to the other, but that does not tell the database which cart is right. The database needs a rule to resolve the conflict. A common rule is to keep the change with the latest timestamp. Here, that keeps one of the shopper’s edits and discards the other.

Strong consistency

In contrast, strong consistency means every read returns the latest write. This is the consistency in the CAP theorem.

It does not mean the replicas must apply a write at the same moment. There is always some delay before every replica applies the write. But a replica that is behind must not answer a read with old data.

One way to provide it is for the primary to wait until all replicas apply the write before it reports success:

sequenceDiagram
    participant A as Shopend API
    participant P as Primary
    participant R1 as Replica 1
    participant R2 as Replica 2
    A->>P: Write
    P->>P: Record write
    P->>R1: Replicate change
    P->>R2: Replicate change
    R1->>R1: Apply change
    R1-->>P: Change applied
    R2->>R2: Apply change
    R2-->>P: Change applied
    P-->>A: Success
    Note over P,R2: Later reads must include this write

Once the API receives success, every replica has the write, so any later read includes it. If any replica cannot respond, the write cannot complete. This is the CP choice from the previous section.

If we need strong consistency for some operations but not all, another way to provide it is to read from the primary for those operations. The primary is already up to date (assuming it is the only primary and all writes go to it). We already do this inside a purchase transaction, and those reads do not wait for any replica.

Read-your-writes

Eventual consistency can confuse users. A shopper adds a mug to their cart and gets a success message. Then they open the cart, the read goes to a replica, and the mug is not there.

Users accept a delay before they see other people’s changes, but they expect to see their own right away. The previous chapter called this guarantee read-your-writes: after you make a change, your later reads include it. In Shopend, we provide it by reading a shopper’s data from the primary for a few seconds after they make a change.

Read-your-writes is stronger than eventual consistency and weaker than strong consistency.