Wide fantasy illustration of a central glowing archive fortress connected to several smaller replicated archive towers by illuminated bridges, with one broken bridge representing network partition and inconsistent data across a distributed system.
The Kingdom in the Clouds

When Two Kings Claim the Throne: Consistency in Distributed Systems

Truth becomes complicated when the kingdom keeps more than one copy of it.

When One Truth Becomes Many: Distributed Systems

The Kingdom Keeps Copies for a Reason: Distributed Data and Replication

In a small application, truth appears simple. A user changes an email address, the application saves it, and the next request retrieves the same value. One database, one authoritative record, and usually one short path between the person making the change and the system storing it.

Distributed systems complicate that arrangement. Data may be copied across regions, databases, caches, queues, search indexes, and services. Replication improves availability and performance, but it also creates a difficult question: when several systems hold a version of the same information, which version should the kingdom believe?

That question is the heart of consistency.

In the previous article, we explored messages and events traveling between distant services. Those messages allowed the kingdom to continue working without forcing every messenger to wait for an immediate reply. However, asynchronous communication also means that information may arrive at different times. One city may know about a change while another still operates with yesterday’s version of events.

This is not automatically a failure. It becomes a failure when the system promises stronger agreement than its architecture can provide.

A distributed system is not merely a single database with extra copies. It is a collection of independent participants that must coordinate across distance, timing, and failure. Once information is replicated, every design decision must answer two related questions: how quickly must every copy agree, and what should the system do while agreement is incomplete?

Answer those questions before choosing a database, replication strategy, or consistency setting. Technology can provide mechanisms, but it cannot decide what truth means for the application.

Replication means maintaining multiple copies of data, usually to improve resilience, reduce latency, or distribute workload. A service in California may read from a database replica in California, while a service in Europe reads from a nearby European replica. If one region becomes unavailable, another may continue serving requests.

The benefits are substantial. Local reads can be faster. Traffic can be distributed across multiple locations. A regional outage does not necessarily become a global outage. Replication helps distributed systems serve large numbers of users across wide geographic areas.

The cost is that the copies may not change at exactly the same moment.

Suppose a customer changes the shipping address for an order. The primary database accepts the update, but replicas receive it a few milliseconds later. During that brief interval, one request may return the new address while another returns the old address. If the application lets the order ship immediately, those two answers are not equally harmless.

The correct response depends on the meaning of the data. A delayed product review may be acceptable. A delayed inventory count may cause overselling. A delayed permission change may become a security problem. A delayed account balance may be unacceptable altogether. The same replication behavior can be reasonable for one feature and dangerous for another.

Consistency is therefore not a single universal setting called safe or unsafe. It is a contract between the system and its users about what they can observe.

The Two Kings at the Throne: Consistency Models

When engineers describe a system as strongly consistent, they usually mean that operations behave as though there were one current copy of the data. A successful write becomes visible in a defined order, and later reads do not casually return older information.

One common strong guarantee is linearizability. Under linearizability, each operation appears to happen at one specific point in time between its beginning and end. The system behaves as if all clients were communicating with one up-to-date authority.

That model is easier to reason about, but it usually requires coordination. Replicas may need to acknowledge writes. Requests may need to travel across regions. A system may have to delay a response until enough participants agree that the operation is durable.

Eventual consistency makes a different promise. If no new updates continue and the system can communicate normally, replicas will eventually converge on the same value. It does not promise that every read immediately returns the newest value.

Eventual consistency can be an excellent choice for social activity feeds, view counts, search indexes, product recommendations, public comments, and cached catalog information. The important word is eventually. It is not a synonym for random, careless, or broken. It is a deliberate tradeoff that lets the system stay responsive while copies catch up.

Between these models are more targeted guarantees. A system might provide read-your-writes consistency, ensuring that a user sees their own update even if other users briefly see an older version. It might provide monotonic reads, preventing one user from seeing a newer version and then later being shown an older one. It might provide causal consistency, preserving the order of related operations.

These narrower guarantees are often more useful than arguing about whether the entire system should be strong or eventual. The best consistency model is the weakest model that safely satisfies the user-facing contract.

When the Roads Close: The CAP Theorem

The CAP theorem describes what happens when a distributed system experiences a network partition, meaning that participating nodes cannot reliably communicate with one another. During that partition, a system cannot guarantee both perfect consistency and continued availability for every operation.

If the system chooses consistency, it may reject or delay requests until it can determine which version of the truth is authoritative. Users may receive errors, but the system avoids accepting conflicting updates.

If the system chooses availability, each reachable part of the kingdom may continue accepting requests. The system remains responsive, but different regions may temporarily accept incompatible versions of the truth.

The network partition is the unavoidable pressure point. Without reliable communication, the system cannot both confirm that every participant agrees and respond immediately to every request.

CAP does not say that distributed systems must always sacrifice consistency. It says a system must make a choice when communication between its parts fails. The design may choose different behavior for different operations, and many systems use consistency levels that vary by request.

Consider a bank account. If two regions independently accept withdrawals from the same account during a network partition, both may believe that enough money exists. When the regions reconnect, the system may discover that the account was overdrawn.

For that operation, availability may be less important than preventing an invalid balance. Rejecting a withdrawal is frustrating. Creating money that does not exist is worse.

Now consider a social media Like button. If a user clicks Like while one region is temporarily disconnected, accepting the action and reconciling the count later may be perfectly reasonable. A temporary count discrepancy does not damage the feature’s underlying meaning.

The architecture should reflect the consequence of being wrong.

The Council Counts Its Voices: Quorum Reads and Writes

Many replicated systems use quorum-based reads and writes to balance consistency, availability, and latency. A quorum is the minimum number of replicas that must participate in an operation.

Suppose a record has five replicas. A system might require three replicas to acknowledge a write before reporting success. It might also require three replicas to respond to a read. Since the read quorum and write quorum overlap, the read has a strong chance of contacting at least one replica that has seen the latest successful write.

The general relationship is:

  • N represents the number of replicas.
  • W represents the number of replicas required to acknowledge a write.
  • R represents the number of replicas required to answer a read.
  • When W + R > N, the read and write quorums overlap.

For five replicas, a configuration such as W = 3 and R = 3 creates overlap. However, quorum overlap does not automatically provide full linearizability. Correctness also depends on version ordering, conflict handling, replica health, clock behavior, and how the database selects the newest response.

A weaker configuration might use W = 2 and R = 1. Writes complete faster and remain available during more failures, but a read can reach a replica that hasn’t received the latest update.

Quorums therefore express a tradeoff rather than a magical guarantee. Larger quorums usually increase confidence that an operation reflects current data, while smaller quorums usually improve availability and reduce latency.

The engineering decision should match the invariant. A profile image may tolerate a stale read. A payment status may require stronger acknowledgment and version validation.

When Both Kings Issue Orders: Conflict Resolution

Replication moves changes between copies. It does not automatically determine what to do when two copies change independently.

Imagine two regional services updating a customer profile while temporarily disconnected. Region A changes the phone number. Region B changes the mailing address. If the system stores the entire profile as one replaceable object, the later update may overwrite the earlier one, losing valid information.

A simple last-write-wins strategy chooses the update with the latest timestamp. This is easy to implement, but it assumes that the latest write should defeat the earlier one. That assumption may be acceptable for a cache. It may be disastrous for a collaborative document or an account record.

Other systems resolve conflicts by field, by version, or through domain-specific rules. A shopping cart may merge independent additions. A calendar may preserve both appointments and flag the collision. A document editor may use an algorithm designed to merge concurrent changes. A financial ledger may refuse to merge conflicting entries and require human review.

The strongest conflict-resolution strategy is often the one that makes invalid states impossible.

Rather than allowing every region to edit every field, the system may assign ownership. One service owns account balances. Another owns catalog descriptions. Another owns delivery status. Other services receive events and maintain projections of the information they need.

Ownership reduces ambiguity because the system does not have to decide which of several competing authorities is correct. It knows where the decision belongs.

When ownership is impossible, the conflict policy must be explicit. Hidden conflict resolution is simply data loss with better public relations.

The Crown’s True Authority: Sources of Truth and Data Ownership

Distributed consistency becomes easier to design when the system separates the source of truth from the views built around it.

A source of truth is the authority that accepts and validates a particular change. A projection is a derived view optimized for reading. The projection may lag behind the source, but the application should understand that delay and communicate it when necessary.

For example, an order service might own the order state. A fulfillment service receives order events and builds a list of shipments. A notification service receives the same events and sends messages to customers. A search service indexes orders for support staff.

Those supporting systems may be eventually consistent. The order service still needs to enforce the rules that prevent an order from being paid twice or shipped before payment. A delayed search result is inconvenient. A duplicate charge is a serious defect.

This separation also helps teams reason about recovery. If a projection becomes corrupted, teams can rebuild it from the authoritative history. If every service independently edits its own copy with no clear authority, recovery becomes an archaeological expedition through contradictory records.

The Scribes Check the Seal: Optimistic Concurrency and Idempotency

One way to protect updates is optimistic concurrency control. The record includes a version number. A client reads version 7, prepares an update, and sends it back expecting version 7 to still be current.

UPDATE accounts
SET mailing_address = :new_address,
    version = version + 1
WHERE account_id = :account_id
  AND version = :expected_version;

If the update affects one row, the change succeeded. If it affects zero rows, another operation changed the record first. The application can reload the current value and ask the user to review the conflict.

The important lesson is not the SQL syntax. It is that the system refuses to silently overwrite a change it has not seen.

This approach works well when conflicts are uncommon, and a human or application can safely retry. It does not solve every distributed data problem. It does, however, turn an invisible race condition into an explicit decision.

Distributed systems also often deliver messages more than once. A consumer may complete its work but fail before acknowledging the message. The broker then resends the message.

If processing is not idempotent, the repeated message may create duplicate work.

def handle_payment_event(event, store):
    if store.has_processed(event["event_id"]):
        return

    store.begin()
    store.record_processed(event["event_id"])
    store.create_receipt(
        payment_id=event["payment_id"],
        amount=event["amount"]
    )
    store.commit()

This example assumes that recording the event and creating the receipt occur within a transaction protected by a unique constraint. Without that protection, two workers could check at the same time, both see that the event is new, and both create a receipt.

Idempotency does not make a system strongly consistent. It makes repeated attempts produce the same meaningful result. That distinction matters. A service can be eventually consistent and still be safe against duplicate delivery.

Reliable distributed systems are often built from several modest guarantees working together:

  • A source of truth owns critical decisions.
  • Version checks prevent silent overwrites.
  • Unique constraints protect identity.
  • Idempotency makes retries safe.
  • Events distribute changes to derived views.
  • Reconciliation detects and repairs drift.
  • User interfaces explain when information may be delayed.

No single mechanism carries the entire burden.

Treaties Between Kingdoms: Distributed Transactions and Sagas

A distributed transaction coordinates changes across multiple services or databases as one logical operation. The classic approach is two-phase commit. A coordinator asks participants whether they can commit, then instructs them to commit after all have agreed.

This can preserve strong consistency, but it introduces coordination overhead and failure complexity. Participants may hold locks while waiting. The coordinator may fail at an awkward moment. Network delays can keep resources unavailable longer than expected.

For some financial or inventory operations, that cost may be justified. For many workflows, a saga is more appropriate. A saga breaks a business transaction into local transactions connected by messages. If a later step fails, the system performs a compensating action.

For example:

  1. Reserve inventory.
  2. Authorize payment.
  3. Create shipment.
  4. Send confirmation.

If shipment creation fails after payment authorization, the system may release the inventory and void the authorization. The operations are not one atomic database transaction, but the workflow has a defined recovery path.

A compensation is not always a perfect reversal. You can’t truly unsend an email. A physical shipment cannot be made to disappear from a customer’s porch. The design must acknowledge those realities and decide which actions are reversible, retryable, or reviewable.

The mature question is not whether a distributed workflow can pretend to be one transaction. It is whether every partial outcome has an understandable and safe next step.

Ruling During Disagreement: Engineering Judgment

When two kings claim the throne, the answer is not always to choose one and destroy the other. Sometimes the kingdom needs one authority. Sometimes it needs temporary local rulers who reconcile later. Sometimes the copies should not be treated as equals at all because one is an authoritative ledger and the others are merely useful maps.

The difficult work is identifying which situation you have.

Before selecting a consistency model, define the user-visible promise. Identify the invariants that must survive concurrency and failure. Decide where ownership lives. Determine how conflicts will be prevented, merged, rejected, or repaired. Then choose the least expensive coordination mechanism that preserves those guarantees.

A distributed system is trustworthy when its behavior remains understandable during disagreement. That may mean rejecting a request, showing a pending state, or accepting a change and reconciling it later. What matters is that the behavior follows a deliberate contract rather than an accidental property of network timing.

The central lesson is simple: consistency is not about forcing every copy to agree instantly. It is about deciding which disagreements are acceptable, which are dangerous, and who has the authority to resolve them.

The Road Ahead: Designing for Network Failure

The Distributed Realm has now gained both messengers and multiple versions of the truth. Messages allow services to work independently, while consistency rules determine how those services behave when information arrives late or copies disagree.

But even the most carefully designed consistency model cannot make the network reliable. Roads fail. Packets disappear. Responses arrive after the caller has stopped waiting. A service may be healthy while the path to it is not.

This Friday, The Unreliable Messenger: Designing for Failure Across the Network will examine timeouts, retries, circuit breakers, idempotency, and partial failure. In The Kingdom in the Clouds, the next lesson is not how to prevent every failure. It is how to build a kingdom that continues serving its people when failure becomes part of ordinary life.

Leave a Reply

Your email address will not be published. Required fields are marked *