Levelwise
English
Distributed systems

Data consistency and the CAP theorem

When data has several copies, the copies are not all the same at every moment. The CAP theorem says that when the network between the copies breaks, you must choose between "data is always correct" and "always answering". For each piece of data, choose the weakest consistency level the business accepts.

Not reviewedWritten with AI helpReading time: 16 minExample of customer address and warehouse stockC# and EF Core code in .NET 10

Author: bezzad

The problem: data that lives in several places

The online shop has grown. The main database (Primary) is tired under the read load. The team adds a read-only copy (Read Replica):

  • All writes go to the main database.
  • All reads are done from the copy.

One day this complaint arrives: “I changed my address and saw a success message. But the page still shows the old address. After a few reloads it was fixed.”

CustomerPrimaryNew address: ShirazReplicaStill: Tehran1. Save new address2. The page reads againDelayed copyThe customer sees the old address
The change reaches the copy with a small delay. The page reads before that.

What happened? Step by step:

  1. The customer pressed save. The address was saved in the main database and the answer “success” came back.
  2. The page reloaded right away and read from the copy.
  3. Copying data to the copy is usually asynchronous. This means the main database does not wait for the copy. This takes from a few milliseconds to a few seconds. This delay is called Replication Lag.
  4. So the page saw the old address. A few seconds later, the change arrived and everything was correct.

There is no bug in the code. This is the nature of data with many copies. The same thing happens between Microservices: when the warehouse service learns about an order change through an event, it is a few seconds behind.

Live example

Change the address once without the option and once with the option:

Live example: change the delivery address

Copying each change from the main database to the copy takes 3 seconds.

TehranMain database
TehranCopy
TehranCustomer page

    Consistency levels

    “Consistency” is not one thing. It has several levels. The stronger the level, the higher the cost:

    StrongEveryone seesthe latest dataRead your writesEach user sees theirown change at onceMonotonic readsData never goesback in timeEventualFinally allbecome the sameSlower, sensitive to breaksFaster, always availableFor each piece of data, choose the weakest level the business accepts
    Consistency levels, from strong to weak.
    1. Strong consistency (Strong or Linearizable). When a write is done, all readers see the new value. It is as if there is only one copy. It is the most expensive level, because each write must wait for the other copies.
    2. Read your own writes (Read-Your-Writes). Each user sees their own change right away. Other users may see it a few seconds later. For the address problem, this is enough.
    3. Monotonic reads (Monotonic Reads). If you saw the new value once, you never see the old value again. Without it, if the first request goes to a fresh copy and the second request goes to a copy that is behind, the data “goes back in time”.
    4. Eventual consistency (Eventual Consistency). If no new write comes, all copies finally become the same. But nobody knows when. It is the cheapest and fastest level.
    The right question: The question is not “which level is better?”. The question is: “how much delay can this specific data accept?”. The like count of a product can accept a few seconds of delay. The account balance before taking money out cannot.

    The CAP theorem

    The CAP theorem is about a system whose data is on several servers. It has three properties:

    • Consistency. Every read sees the latest write. Here it means strong consistency, not the letter C in ACID.
    • Availability. Every healthy server answers every request.
    • Partition Tolerance. The system works even when the connection between servers is broken.

    The theorem says: when the network between servers breaks, you cannot have both full consistency and full availability. You must choose one.

    Tehran data centerLast item: soldStock: zeroTabriz data centerDoes not know yetStock: oneConnection is brokenA customer in Tabriz wants the same item. What should Tabriz do?Choose consistencyReject: "Not now, try again later"No wrong data, but also no answerChoose availabilityAccept: sell with maybe old dataWe answer, but may need to compensate later
    The connection between two data centers is broken. Each choice has a cost.

    Why? Step by step:

    1. The shop has two data centers: Tehran and Tabriz. Both have a copy of the stock.
    2. The connection between them breaks.
    3. In Tehran, the last item is sold. Tabriz does not know about it.
    4. Another customer in Tabriz wants the same item. Tabriz has two ways:
      • Reject until the connection comes back. It does not give wrong data, but it is not available. This is called CP.
      • Accept with data that may be old. It is available, but one item may be sold twice. This is called AP.

    Three common misunderstandings

    1. “Choose two of three” is not exact. A network break in a distributed system is not a choice; it happens. So the real choice is only this: during a break, consistency or availability? A “CA” system in practice means a system that runs on one server.
    2. This choice is not for the whole system. You can decide separately for each piece of data and each action. The shopping cart can be AP and payment can be CP.
    3. Network breaks are rare, but delay is always there. That is why we also have a fuller model called PACELC: if there is a break (P), choose between A and C. Else (Else), choose between low latency (L) and consistency (C). This means that even on normal days, stronger consistency means a slower answer.

    What should we choose in the business?

    Data or action Sensible choice Why?
    Account balance before a withdrawal Strong consistency Taking out more than the balance is a real loss.
    Booking the last seat or the last item Strong consistency, or sell and compensate It depends: is rejecting the customer worse, or saying sorry?
    User address and profile Read your own writes The user must see their own change. Others are in no hurry.
    Shopping cart Availability We must never say “you cannot add to the cart right now”.
    View and like counts Eventual consistency A few seconds of delay does not matter to anyone.
    Manager report page Eventual consistency Data from a few minutes ago is enough.
    A third way: Sometimes the best thing is to let the mismatch happen and then compensate for it. For example, if an item was sold twice, give the second customer a discount code and an apology. This is a business decision and must be made with the product manager.

    Eventual consistency between services

    In Microservices, services stay in sync with events. So between them we always have eventual consistency. Three points for living with it:

    1. An honest user interface. After an order is placed, write “processing”, not “confirmed”. Keep it that way until all services have done their work.
    2. Duplicate events. Each event may arrive twice. So the receiver must be Idempotent.
    3. Events in the wrong order. The event “address changed to Shiraz” may arrive after the newer event “address changed to Tabriz”. You also cannot trust the server clocks, because clocks are not exactly the same. The right way: a version number that grows in only one place (the service that owns the data). The receiver accepts only a newer version.

    Code

    Read your own writes with EF Core

    We have two DbContexts: one for the main database and one for the copy. After each write, we set a short-lived cookie. While this cookie exists, we read from the main database:

    builder.Services.AddDbContext<ShopDbContext>(o => o.UseNpgsql(primaryConnection));
    builder.Services.AddDbContext<ShopReadDbContext>(o => o.UseNpgsql(replicaConnection));
    
    app.MapPut("/customers/{id:guid}/address", async (
        Guid id, AddressDto dto, ShopDbContext db, HttpContext http, CancellationToken ct) =>
    {
        var customer = await db.Customers.FindAsync([id], ct);
        if (customer is null) return Results.NotFound();
    
        customer.ChangeAddress(dto.City, dto.Street);
        await db.SaveChangesAsync(ct);
    
        // For the next 10 seconds, this user reads from the primary.
        http.Response.Cookies.Append("recent-write", "1",
            new CookieOptions { HttpOnly = true, MaxAge = TimeSpan.FromSeconds(10) });
        return Results.NoContent();
    });
    
    app.MapGet("/customers/{id:guid}", async (
        Guid id, HttpContext http, ShopDbContext primary, ShopReadDbContext replica, CancellationToken ct) =>
    {
        IQueryable<Customer> customers = http.Request.Cookies.ContainsKey("recent-write")
            ? primary.Customers
            : replica.Customers;
    
        var customer = await customers.AsNoTracking().SingleOrDefaultAsync(c => c.Id == id, ct);
        return customer is null ? Results.NotFound() : Results.Ok(customer);
    });

    This is simple and keeps most of the read load on the copy. The 10 seconds must be more than the usual copy delay. So monitor the copy delay and set an alert for it.

    Accept only a newer version

    The order service has a local copy of the customer address. It is updated by an event:

    public sealed record CustomerAddressChanged(Guid CustomerId, string City, long Version);
    
    public sealed class CustomerAddressChangedHandler(OrderDbContext db)
    {
        public Task HandleAsync(CustomerAddressChanged e, CancellationToken ct) =>
            // Old or duplicate events match no row, so they change nothing.
            db.CustomerCopies
                .Where(c => c.Id == e.CustomerId && c.Version < e.Version)
                .ExecuteUpdateAsync(s => s
                    .SetProperty(c => c.City, e.City)
                    .SetProperty(c => c.Version, e.Version), ct);
    }

    This is a conditional update command. If the event is older or a duplicate, no row matches the condition and nothing changes. So this code solves both the wrong order and the duplicate message.

    Important rules

    1. Decide separately for each piece of data. One consistency level for the whole system is either expensive or dangerous.
    2. Make important decisions on fresh data. A balance check before a withdrawal, a permission check, and any read that you write after, must come from the main database.
    3. The user must see their own change. This is the least a user expects.
    4. Understand order with a version number, not with a clock. The clocks of different servers are not exactly the same.
    5. Measure the copy delay. Without a number, you do not know if eventual consistency means one second or one hour.
    6. Design an honest user interface. “Processing” is better than “confirmed” that gets canceled later.

    Common mistakes

    Mistake Result The right way
    All reads from the copy The user does not see their own change. Read your own writes.
    Checking the account balance from the copy Taking out more than the balance. Important decisions from the main database.
    Sorting events by server clocks Wrong order and wrong reports. A version number, or the order inside one Partition in Kafka.
    Thinking “CAP means two of three” Design based on a wrong idea. A network break is always possible. The choice is between C and A.
    Strong consistency for everything A slow system that is sensitive to breaks. The weakest level the business accepts.
    Showing “confirmed” before the services are in sync The customer thinks the job is done. A “processing” status.

    Eventual consistency is fine

    • Display data like view counts and product suggestions.
    • Data copies between Microservices.
    • Reports and search.

    Eventual consistency is dangerous

    • Money, account balance and payment.
    • Counters that must not go over a limit, like the last item.
    • Permissions, for example after a user’s access is taken away.

    Summary in six lines

    1. When data has several copies, the copies become the same with a delay.
    2. Consistency levels go from strong to eventual. The stronger, the slower and more expensive.
    3. The CAP theorem says that during a network break, you must choose between consistency and availability.
    4. Even without a break, stronger consistency means more delay (PACELC).
    5. Decide separately for each piece of data, and let each user see their own change right away.
    6. Between services, live with eventual consistency using version numbers and Idempotent receivers.