Replicate data across machines and you inherit a question a single-server database never had to answer: when two nodes disagree about a value, which one is right? Every distributed system answers with a consistency model, and the answer determines whether your application can trust what it reads. This page defines data consistency, compares the three models you will choose between, and shows the mechanisms 티DB uses to keep strong consistency without giving up availability.

What Is Data Consistency in Distributed Systems?

Data consistency in a distributed system is the guarantee that every node presents the same view of the data, so a read returns a value the system agrees on rather than whichever copy the request reached. The guarantee is set by the system’s consistency model, which defines the recency and ordering an application can rely on. Strong consistency means every read returns the most recent committed write; weaker models allow reads to lag or to observe operations in different orders.

Three models cover most decisions. Strong consistency guarantees that a read returns the latest committed write, and that all clients observe operations in one order. Eventual consistency guarantees only that replicas converge once writes stop, so a read may return a stale value meanwhile. Causal consistency sits between them: causally related operations are observed in order, while unrelated concurrent ones may differ per replica.

Why Consistency Is Hard to Guarantee at Scale

Three physical realities make this difficult. Network partitions cut nodes off with no way to tell a slow peer from a dead one. Replication latency means a write acknowledged on one node has not reached the others. Node failures remove replicas mid-operation, sometimes after a write is accepted but before it is durable elsewhere.

The CAP theorem states the resulting constraint: during a network partition, a system can preserve consistency or availability, not both. Consider a transfer that debits one account and credits another. If the system stays available during a partition and accepts writes on both sides, both can debit the same balance, and reconciliation happens after the money is gone. If it preserves consistency, the minority side refuses writes and the transfer fails cleanly, which is recoverable. For financial data, the second failure is the one you want.

Consistency Models Compared

ModelGuaranteeLatency costExample use case
Strong (linearizable)Read returns the latest committed write; one global orderHighest: coordination on every write, floor set by replica distancePayment ledgers, inventory, order books
CausalCausally related operations ordered; concurrent ones may vary per replicaModerate: only related operations need orderingMessaging, comment threads, collaborative editing
EventualReplicas converge once writes stopLowest: no coordination before acknowledgingView counters, DNS, cache invalidation

Stronger models put coordination cost in the database; weaker ones put reconciliation code in your application. For the client-centric guarantees a user actually notices, such as read-your-writes, see consistency models in distributed systems.

How TiDB Solves the Consistency vs. Availability Trade-off

TiDB chooses consistency, and the mechanism is specific enough to be worth naming precisely.

Data is split into Regions, contiguous ranges of key-value pairs. Each Region’s replicas form their own independent Raft group with their own leader, so a cluster holding thousands of Regions runs thousands of Raft groups concurrently. That is what Multi-Raft describes. This matters for availability: a write commits once a majority of that Region’s replicas have persisted it, so a node failure cannot lose a committed write, and leadership spreads across the cluster rather than concentrating in one coordinator. If a leader fails, the remaining replicas elect a new one automatically without losing a committed write. See TiDB’s architecture and the wider distributed SQL architecture overview.

Raft covers replication within a Region. Transactions spanning Regions need more, and TiDB uses a Percolator-style two-phase commit over multi-version concurrency control. One key acts as the primary and the rest as secondaries; committing the primary key is the atomic decision point, so a transaction touching many Regions commits everywhere or nowhere. MVCC lets reads proceed against a stable snapshot without blocking on in-flight writes, giving Snapshot Isolation by default with no configuration required.

Data Consistency for AI Agent Workloads

Agent workloads make consistency a correctness problem rather than a user-experience one. An agent reads state, decides, writes, then reads again on the next step. If a read returns a stale value, the agent may act on state it has already changed, and unlike a human user it will not notice and retry. Eventual consistency turns every read into a decision on possibly outdated information.

Kimi’s agent hosting platform shows what predictable reads buy at scale. Kimi, the AI product from Moonshot AI, runs a K2.6 agent that turns a plain-language request into a deployed and hosted web application, frontend and backend included. It runs tens of millions of concurrent tenant sites on TiDB Cloud, each agent task receiving a fully prepared database in under one second. The detail that matters here: the agent generates application code with no retry or polling logic for database availability. Every extra boundary an agent reasons across adds a class of failure it must handle, and predictable reads remove one of them.

Test the guarantees against your own workload. Start with TiDB Cloud Starter, free.

Frequently Asked Questions

What is the difference between strong and eventual consistency?

Strong consistency guarantees that a read returns the most recently committed write and that all clients observe operations in the same order, as though only one copy existed. Eventual consistency guarantees only that replicas converge once writes stop, so a read may return a stale value and two simultaneous reads against different replicas may disagree. Strong consistency puts coordination cost in the database; eventual consistency pushes conflict resolution into the application.

Does TiDB guarantee strong consistency?

Yes, by default and without configuration. Each Region’s replicas form a Raft group that commits a write only after a majority have persisted it, and reads are served from the Region leader. Transactions spanning Regions use a Percolator-style two-phase commit over MVCC, so they commit everywhere or nowhere. The isolation level is Snapshot Isolation, exposed as REPEATABLE-READ for MySQL compatibility.

How does the CAP theorem apply to distributed SQL databases?

The CAP theorem says that during a network partition a system must choose between consistency and availability. Distributed SQL databases including TiDB choose consistency: the minority side of a partition stops accepting writes rather than accepting writes it cannot reconcile. The majority side keeps serving, so the effect is not total unavailability but reduced availability confined to the partitioned minority. CAP constrains behavior during failures, not normal operation.

Can a distributed database be both consistent and highly available?

In everyday operation, yes. High availability comes from replication and automatic failover, and TiDB delivers both while remaining strongly consistent, since a majority of replicas can commit and serve reads even when nodes fail. What CAP rules out is preserving both during a network partition. A system can survive many node failures with no consistency loss and still refuse writes on the minority side of a partition, which is why availability targets should be stated as a tolerated number of failures rather than an absolute.


Last updated 11월 30, 2024

💬 Let’s Build Better Experiences — Together

Join our Discord to ask questions, share wins, and shape what’s next.

Join Now