TiDB looks like MySQL until you ask where the query runs
· Andres Kepler · TiDB Engineering
A MySQL client connects to TiDB and nothing looks different. Behind that connection are TiDB servers, TiKV nodes and PD, and the surprises come from how they interact.
Point a MySQL client at TiDB and it connects as if TiDB were MySQL. The protocol is the same. The database behind it is not one server but several kinds of process, each scaled on its own.
Most production surprises after a move from MySQL come from how those layers interact, not from SQL itself: hotspots, latency, plan changes. So when a cluster misbehaves, I start with the query plan.
Three processes, three jobs
The TiDB server is the SQL layer. It parses the query, optimizes it with a cost-based optimizer and executes it. It stores no data, so any TiDB server can take any connection. They sit behind a load balancer or ProxySQL.
TiKV is the storage layer. Rows and indexes are sorted key-value pairs, split into ranges called Regions. Each Region is replicated with Raft, three replicas by default, and one replica is the leader.
PD, the Placement Driver, is the coordinator. It keeps the cluster metadata, issues the timestamp every transaction needs, and schedules Regions: splitting, merging, and moving leaders and replicas to balance load.
Follow one query
A query starts on a TiDB server. The server gets a timestamp from PD and builds a plan from the table statistics. Then it pushes as much work as it can down to the TiKV nodes that hold the relevant Regions: filters, aggregations, limits.
EXPLAIN shows where each step runs. Work on TiKV appears as cop[tikv]. 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. The same query can be much faster if its work spreads across many nodes. Same SQL, different cost, depending on where the work runs.
The surprises follow from the layers
A sequential AUTO_INCREMENT primary key sends every insert to the same Region on the same TiKV node. Many machines, one of them busy. AUTO_RANDOM spreads the writes.
The optimizer is cost-based, so stale statistics can make it pick a different index or join overnight. The SQL did not change. Stale statistics were enough.
Capacity runs out differently in each layer: CPU on the TiDB servers, disk and I/O on TiKV, scheduling and timestamp pressure on PD. Each layer scales by adding nodes, and scaling the wrong one buys nothing. More TiDB servers will not help a cluster that is short of disk I/O.
Compatibility belongs at the start, not the end. TiDB is MySQL-compatible, not MySQL. Stored procedures, triggers and events are not supported, and some behaviours differ. A migration starts by finding out what the application relies on.
The full guide, including transactions and recovery, is at https://moitoi.tech/tidb/how-tidb-works.
The detail: moitoi.tech/tidb/how-tidb-works
