MoiToi.TECHTiDB EngineeringGuide
How TiDB works
Updated · Andres Kepler
In short
TiDB splits a MySQL-compatible database into three layers: stateless TiDB servers that speak SQL, TiKV nodes that store the data as replicated key-value Regions, and PD, which hands out transaction timestamps and decides where data lives. Each layer scales by adding nodes. Most production surprises — hotspots, latency, plan changes — come from how those layers interact, not from SQL itself.
01
The three layers
A TiDB cluster is not one server. It is several kinds of process, each scaled on its own.
- — TiDB server — the SQL layer. It speaks the MySQL protocol, parses and optimizes queries with a cost-based optimizer, and executes them. It stores no data, so any TiDB server can take any connection; they sit behind a load balancer or ProxySQL.
- — TiKV — the storage layer. Rows and indexes are stored as sorted key-value pairs and split into ranges called Regions. Each Region is replicated (three replicas by default) with the Raft consensus protocol; one replica, the leader, serves reads and writes.
- — PD (Placement Driver) — the coordinator. It keeps cluster metadata, issues the timestamp (TSO) every transaction needs, and schedules Regions: splitting, merging and moving leaders and replicas to balance load.
- — TiFlash — optional columnar replicas of chosen tables, kept in sync as Raft learners, so analytical queries can run without loading the row store.
02
What happens when a query runs
The client connects to a TiDB server as if it were MySQL. TiDB gets a timestamp from PD, builds a plan from the table statistics, and pushes as much work as possible — filters, aggregations, limits — down to the TiKV (or TiFlash) nodes that hold the relevant Regions. In EXPLAIN output that work appears as cop[tikv] or cop[tiflash]; the part that runs on the TiDB server appears as root.
Every step that crosses the network adds latency. A query that is fast on a single MySQL primary can be slower on TiDB if it touches many Regions or cannot push work down — and much faster if it can be spread across many nodes.
03
Transactions and MVCC
TiDB runs distributed transactions with two-phase commit, based on Google's Percolator model, and keeps multiple versions of each row (MVCC) for snapshot reads. Pessimistic locking is the default, which behaves much closer to MySQL InnoDB than the original optimistic mode.
Old row versions are removed by garbage collection once they are older than tidb_gc_life_time (10 minutes by default). Long-running queries, exports and replication that need older snapshots depend on that window.
04
What this means in production
The architecture explains most of the problems teams run into after moving from MySQL:
- — Write hotspots — a sequential AUTO_INCREMENT primary key sends every insert to the same Region and the same TiKV node. AUTO_RANDOM or SHARD_ROW_ID_BITS spread the writes.
- — Plan changes — the optimizer is cost-based, so stale statistics can make it choose a different index or join overnight. See TiDB query plans and statistics.
- — Compatibility — TiDB is MySQL-compatible, not MySQL. Stored procedures, triggers and events are not supported, and some behaviours differ; a migration starts by finding what the application relies on.
- — Capacity — each layer runs out differently: CPU on TiDB servers, disk and I/O on TiKV, scheduling and TSO pressure on PD. Scaling the wrong layer buys nothing.
- — Recovery — data is spread over many nodes, so backup and restore are cluster operations with their own tooling. See TiDB backup and PITR.
Next step