Why Sharding May Become Obsolete: TiDB's NewSQL Architecture Explained
This article explains how TiDB's NewSQL architecture eliminates manual sharding by providing horizontal scaling, ACID transactions, and MySQL compatibility through a distributed architecture with TiKV storage, PD scheduling, and TiFlash for HTAP workloads.
Introduction: The Problem with Traditional SQL and NoSQL
Traditional single-node MySQL hits a scaling ceiling; upgrading hardware only delays the limit. Sharding across multiple MySQL instances introduces middleware complexity, especially for cross-shard joins and distributed transactions, often pushing logic into the application layer.
NoSQL databases sacrifice strong consistency and SQL support for scalability and flexibility, making them unsuitable for enterprise applications requiring ACID guarantees.
NewSQL: Combining SQL and NoSQL Strengths
NewSQL retains the relational model and SQL while adding NoSQL-like horizontal scalability. It provides SQL support for complex queries, full ACID transactions with isolation levels, elastic scaling transparent to applications, and automatic high availability with failover.
TiDB Origins: From Google Spanner/F1 Papers
PingCAP co-founder and CTO Huang Dongxu was inspired by Google's 2012 Spanner/F1 papers describing a globally distributed relational database with elastic scaling and strong consistency. TiDB was built on those principles.
TiDB Core Features
Horizontal elastic scaling: add nodes to scale throughput or storage on demand.
Distributed ACID transactions: 100% standard ACID support across shards.
Financial-grade high availability: Multi-Raft consensus ensures strong consistency; auto-failover without manual intervention.
Real-time HTAP: row-based TiKV and columnar TiFlash engines replicate data in real time via Multi-Raft Learner, enabling OLTP and OLAP on a single dataset without ETL.
Cloud-native: designed for Kubernetes, runs on public, private, or hybrid clouds.
High MySQL compatibility: MySQL 5.7 protocol and syntax; most applications migrate with zero or minimal code changes.
TiDB Architecture
TiDB Server
Stateless SQL layer that parses, optimizes, and executes queries. It requests data locations from PD, interacts with TiKV/TiFlash, and returns results. Horizontally scalable behind a load balancer (LVS, HAProxy, F5).
PD (Placement Driver) Server
Cluster management module with three roles: stores cluster metadata (key-to-TiKV mapping), schedules load balancing (Region migration, Raft leader transfer), and allocates globally unique monotonic transaction IDs. PD uses Raft for safety; deploy an odd number of nodes (recommended three).
TiKV Server
Distributed transactional key-value storage engine. Data is divided into Regions (key ranges, default ~144 MB). Each Region has multiple replicas (default three) forming a Raft group. PD schedules Regions for balanced distribution. TiKV handles Region splitting/merging automatically.
TiFlash and TiSpark
TiFlash stores data in columnar format, replicating from TiKV via Multi-Raft Learner for strong consistency. It can be deployed on separate machines for resource isolation. TiSpark runs Spark SQL directly on TiKV, integrating with the big-data ecosystem for complex OLAP.
TiKV Internal Architecture
Region Splitting and Merging
When a Region exceeds 144 MB, TiKV splits it into two or more Regions. When deletions shrink adjacent Regions, TiKV merges them. This keeps Region sizes uniform for effective PD scheduling.
Region Scheduling via Raft
Writes go to the Region leader and must succeed on a majority of replicas (default 3 replicas → 2 acknowledgments). PD schedules replica movement by adding a Learner on the target node, promoting it to Follower once caught up, then removing the source Follower. Leader transfer follows a similar process with an explicit leader election step.
Distributed Transactions with Two-Phase Commit
TiKV supports multi-key transactions across Regions using a two-phase commit protocol, guaranteeing ACID without requiring the client to know data placement.
High Availability Design
TiDB High Availability
Stateless; deploy at least two instances behind a load balancer. Single instance failure affects only active sessions; clients reconnect to another instance. Recovery by restarting or adding a new instance.
PD High Availability
PD cluster uses Raft. Non-leader failure has no impact; leader failure triggers a new election (~3 seconds). Deploy at least three PD nodes.
TiKV High Availability
Raft replication with configurable replica count (default three). Node failure affects Regions hosted on that node. Leader loss triggers re-election; Follower loss is transparent. If a node remains down for 10 minutes (default), PD migrates its Regions to healthy nodes.
Key Use Cases
MySQL Sharding Consolidation
TiDB acts as a MySQL-compatible replica (via Syncer) for existing sharded MySQL clusters, enabling real-time cross-shard SQL queries. Huang Dongxu noted: "Past databases were one-master-many-slave; with TiDB you can do many-master-one-slave."
Direct MySQL Replacement
For systems that never sharded, TiDB replaces MySQL directly. All distributed logic stays in the database layer; no application-level sharding code. Horizontal scaling handles growth with stable latency.
Data Warehouse / OLAP
TiDB 2.0 runs TPC-H analytical queries mostly under 10 seconds. For heavier OLAP, TiSpark allows Spark SQL on TiKV in real time.
Embedded Key-Value Store
TiKV can serve as an HBase replacement with cross-row ACID transactions. Two APIs: ACID Transaction API (multi-row) and Raw API (single-row, higher performance). Some users implement Redis protocol on TiKV for high-capacity, latency-tolerant caches.
MySQL Compatibility and Differences
Supported Features
TiDB supports MySQL 5.7 protocol and most syntax. Existing connectors (PHPMyAdmin, Navicat, MySQL Workbench) and tools (mysqldump, Mydumper/myloader) work directly. Most applications migrate without code changes.
Unsupported MySQL Features
Stored procedures and functions
Triggers
Events
User-defined functions
Foreign key constraints
Temporary tables
Full-text/spatial functions and indexes
Non-ascii/latin1/binary/utf8/utf8mb4 character sets (only utf8mb4)
SYS schema
MySQL optimizer trace
XML functions
X-Protocol
Savepoints
Column-level privileges
XA syntax (internal two-phase commit not exposed via SQL)
CREATE TABLE ... AS SELECT CHECK TABLE,
CHECKSUM TABLE GET_LOCK/
RELEASE_LOCKAuto-Increment Behavior Differences
TiDB auto-increment guarantees uniqueness and monotonic increase per TiDB server, but not across servers nor continuity. Mixing default and explicit values may cause Duplicated Error. The system variable tidb_allow_remove_auto_inc controls removal of AUTO_INCREMENT attribute via ALTER TABLE ... MODIFY or CHANGE; once removed it cannot be re-added.
SELECT Statement Limitations
No SELECT ... INTO @variable.
No SELECT ... GROUP BY ... WITH ROLLUP. SELECT ... GROUP BY expr returns no guaranteed order (matches MySQL 8.0), whereas MySQL 5.7 implicitly orders by the grouping expression.
View Restrictions
Views are read-only; UPDATE, INSERT, DELETE on views are not supported.
Default Configuration Differences
Character Set
TiDB defaults to utf8mb4; MySQL 5.7 defaults to latin1; MySQL 8.0 defaults to utf8mb4.
Collation
TiDB utf8mb4 default collation is utf8mb4_bin (binary). MySQL 5.7 uses utf8mb4_general_ci; MySQL 8.0 uses utf8mb4_0900_ai_ci.
lower_case_table_names
TiDB only supports value 2 (store as given case, compare case-insensitively). MySQL defaults: Linux 0 (case-sensitive), Windows 1 (store lowercase, compare case-insensitive), macOS 2.
timestamp Auto-Update
TiDB forces explicit_defaults_for_timestamp=ON (timestamp columns do not auto-update on row change). MySQL 5.7 defaults to OFF (auto-update); MySQL 8.0 defaults to ON.
Foreign Key Support
TiDB forces foreign_key_checks=OFF; MySQL 5.7 defaults to ON.
OLTP vs OLAP Comparison
The article provides detailed comparison tables. Key contrasts:
Real-time requirements: OLTP demands immediate transaction processing; OLAP tolerates daily batch updates.
Data volume: OLTP handles tens of records per transaction; OLAP scans millions for aggregations.
User orientation: OLTP serves customers/front-line staff; OLAP serves decision-makers.
Primary operations: OLTP: insert/update/delete; OLAP: complex queries.
Schema design: OLTP uses ER modeling; OLAP uses star/snowflake schemas.
Data characteristics: OLTP: current, detailed, 2D; OLAP: historical, aggregated, multidimensional.
Access pattern: OLTP reads/writes tens of rows; OLAP reads millions.
Work unit: OLTP: simple transactions; OLAP: complex queries.
Concurrent users: OLTP: thousands; OLAP: hundreds.
Database size: OLTP: 100 MB–GB; OLAP: 100 GB–TB.
Signed-in readers can open the original source through BestHub's protected redirect.
This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactand we will review it promptly.
Architect's Guide
Dedicated to sharing programmer-architect skills—Java backend, system, microservice, and distributed architectures—to help you become a senior architect.
How this landed with the community
Was this worth your time?
0 Comments
Thoughtful readers leave field notes, pushback, and hard-won operational detail here.
