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.
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.”
What happened? Step by step:
- The customer pressed save. The address was saved in the main database and the answer “success” came back.
- The page reloaded right away and read from the copy.
- 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.
- 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:
Copying each change from the main database to the copy takes 3 seconds.
Consistency levels
“Consistency” is not one thing. It has several levels. The stronger the level, the higher the cost:
- 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.
- 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.
- 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”.
- 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 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.
Why? Step by step:
- The shop has two data centers: Tehran and Tabriz. Both have a copy of the stock.
- The connection between them breaks.
- In Tehran, the last item is sold. Tabriz does not know about it.
- 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
- “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.
- 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.
- 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. |
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:
- An honest user interface. After an order is placed, write “processing”, not “confirmed”. Keep it that way until all services have done their work.
- Duplicate events. Each event may arrive twice. So the receiver must be Idempotent.
- 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
- Decide separately for each piece of data. One consistency level for the whole system is either expensive or dangerous.
- 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.
- The user must see their own change. This is the least a user expects.
- Understand order with a version number, not with a clock. The clocks of different servers are not exactly the same.
- Measure the copy delay. Without a number, you do not know if eventual consistency means one second or one hour.
- 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
- When data has several copies, the copies become the same with a delay.
- Consistency levels go from strong to eventual. The stronger, the slower and more expensive.
- The CAP theorem says that during a network break, you must choose between consistency and availability.
- Even without a break, stronger consistency means more delay (PACELC).
- Decide separately for each piece of data, and let each user see their own change right away.
- Between services, live with eventual consistency using version numbers and Idempotent receivers.