Databases 26 min read

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.

Architect's Guide
Architect's Guide
Architect's Guide
Why Sharding May Become Obsolete: TiDB's NewSQL Architecture Explained

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_LOCK

Auto-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.

Original Source

Signed-in readers can open the original source through BestHub's protected redirect.

Sign in to view source
Republication Notice

This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactadmin@besthub.devand we will review it promptly.

distributed databaseTiDBHTAPdatabase shardingNewSQLMySQL compatibilityRaft consensusTiKV
Architect's Guide
Written by

Architect's Guide

Dedicated to sharing programmer-architect skills—Java backend, system, microservice, and distributed architectures—to help you become a senior architect.

0 followers
Reader feedback

How this landed with the community

Sign in to like

Rate this article

Was this worth your time?

Sign in to rate
Discussion

0 Comments

Thoughtful readers leave field notes, pushback, and hard-won operational detail here.