NOTE
Database Sharding and Table Sharding
1. Why Database Sharding and Table Sharding Are Needed Distributed System Partitioning.md 2. What Database Sharding and Table Sharding Are Database sharding splits one original database into multiple databases; table sharding splits one original table into multiple tables. Both database sharding and table sharding can use horizontal or vertical splitting. 2.1. Vertical Splitting 2.1.1. Vertical Table Splitting Each table has a different structure and different data. Split the fields of one table into multiple tables according to usage frequency and whether they are large fields (one-to-one relationship). For example, split a product information table into a product basic information table and a product description table. This operation is generally completed during the initial design. 2.1.2. Vertical Database Splitting Each database has a different structure and different data.
This is a historical learning note and may contain outdated or incomplete understanding.
1. Why Database Sharding and Table Sharding Are Needed
Distributed System Partitioning.md
2. What Database Sharding and Table Sharding Are
Database sharding splits one original database into multiple databases; table sharding splits one original table into multiple tables. Both database sharding and table sharding can use horizontal or vertical splitting.
2.1. Vertical Splitting
2.1.1. Vertical Table Splitting
Each table has a different structure, and each table has different data. Split the fields of one table into multiple tables according to usage frequency and whether they are large fields (one-to-one relationship). For example, split a product information table into a product basic information table and a product description table. This operation is generally completed during the initial design.
2.1.2. Vertical Database Splitting
Each database has a different structure, and each database has different data. Split one database into multiple databases according to the coupling degree of business logic and deploy them on different machines. This is essentially the same logic as splitting microservices by business: each service corresponds to one database. For example, a product information database can be split into a product database and a shop database. This operation is generally completed during the initial design.
2.2. Horizontal Splitting
2.2.1. Horizontal Table Splitting
Each table has the same structure, but each table has different data.
When the number of rows in a single table is very large, split the rows among different tables. For example, if a product table has 500,000 rows, split it into 5 tables, each with 100,000 rows. When the data volume is medium, this can slightly improve performance. When the data volume becomes large, horizontal database sharding is still needed.
2.2.2. Horizontal Database Splitting
Each database has the same structure, but each database has different data. When the number of rows in a single table is very large, split the rows among different databases. For example, if a product database has 5 million rows, split it into 50 databases, each with 100,000 rows. This operation is used only when there are still performance problems after database index optimization, adding cache, and read/write separation.
3. Problems Brought by Database Sharding and Table Sharding
3.1. Primary-Key ID Uniqueness
3.2. Transaction Consistency
3.3. Cross-Node Sorting, Aggregation, and Pagination
- Query concurrently and wrap the results in the code layer. For example, for sorting, send requests to all nodes, take the first N rows from each, aggregate them, sort them, and then take the first N rows. See Elasticsearch CRUD Process.md.
3.4. Cross-Node Join Problems
- Use single-table queries and wrap the results in the code layer.
- Use a wide table, making writes heavier and reads lighter.
- Use a search engine.
4. Database and Table Sharding Middleware
4.1. Classification
- Mainly divided into client and proxy types
- sharding-jdbc (client)
- mycat (proxy)
4.2. Sharding JDBC
4.2.1. What It Is
A client-side database and table sharding middleware. It can be understood as an enhanced JDBC driver. Its core functions include data sharding and read/write separation.
4.2.2. Unsupported Items
A subquery must have a sharding-column.
Batch insertion is not supported.
Aggregate functions must use aliases.
SQL cannot contain parentheses, for example in limit statements.
5. Horizontal Splitting Schemes
5.1. Split Data
5.1.1. Choosing the Sharding Key
Primary-key ID, business ID, or time, balancing query efficiency and even distribution.
5.1.2. Sharding Strategy
5.1.2.1. Separate Mapping Table
- Separate mapping table
- Advantage: the mapping algorithm between IDs and databases can be changed arbitrarily.
- Disadvantages:
- If the mapping database itself becomes too large, it also needs database and table sharding.
- Writes need to write both the business database and the mapping database, which is a distributed transaction problem. Eventual consistency can be guaranteed by scheduled comparison tasks.
5.1.2.2. Range Splitting / Sequential Splitting

- Advantage: the size of a single table is controllable, and horizontal scaling is natural.
- Disadvantage: there can be a hotspot data problem. Newly added data is concentrated on one node, causing a write bottleneck.
- Suitable for time-based splitting in hot/cold data separation scenarios.
5.1.2.3. Hash Splitting

- Advantage: there is basically no hotspot problem.
- Disadvantage: data needs to be migrated during scaling.
- Consistent Hashing must be used. Either use hash slots like Redis, meaning the number of databases is planned in advance, or use a hash ring.
5.2. Request Processing
5.2.1. Routing Component
Either a client or a proxy.
5.2.2. Primary Index
Distributed System Partitioning - Request Processing.md
5.2.3. Secondary Index
- If querying by a non-sharding key, requests need to be sent to all nodes and the results aggregated at the end. But this is inefficient.
5.2.3.1. Redundant Two Copies of the Data
Keep two redundant copies of the data: one sharded by ID1 and the other sharded by ID2. For writes, the business layer can dual-write, or it can single-write plus listen to the binlog.
5.2.3.2. Unified Dimension
For example, use ID1 as the first several digits of ID2, so ID2 contains information about ID1. Querying by ID1 works directly; when querying by ID2, first extract ID1 and then perform the query.
5.3. Partition Assignment
Distributed System Partitioning - Partition Assignment.md This is the data migration needed during scaling in or out and needs to be handled manually.
5.3.1. Downtime
During the early morning when there are no users, stop the service, pull the data from the old single database, write it again into the new sharded databases, and after that is complete, modify the code to connect to the sharded databases and resume service.
5.3.2. Dual Writing Without Downtime
Modify the code so that writes go to both the old single database and the new sharded databases. For old data in the old database, use a script to pull it out and write it again into the sharded databases.
6. How Much Data in a Single Table Before Considering Database and Table Sharding
20 million rows, because the corresponding tree height is 3 and the number of I/O operations is approximately 3.
For the primary-key index B+ tree of the MySQL InnoDB storage engine, assume the page size is 16 KB, the primary key is of type bigint, and each row is 1 KB. How many records can three levels hold?
With a page size of 16 KB, an 8 B key, and a 6 B pointer, one page can store 16*1024/(8+6)=1170 keys.
With a page size of 16 KB and a 1 KB value, one page can store 16/1=16 values.
Assume the tree height is 3:
Level 0: root level. Resident in memory.
Level 1: key level. Can store 1170 keys, corresponding to 1170 branches.
Level 2: key level. Can store 1170*1170 keys, corresponding to 1368900 branches.
Level 3: value level. Can store 1170*1170*16 values, corresponding to 21902400 (about 20 million) rows.
.
The first two key levels occupy about (1170+1368900)*16/1024/1024=21MB, so indexes at this level can be cached in memory. One level higher would occupy (1170+1368900+1601613000)*16/1024/1024=24GB, which cannot be cached in memory. The final value level occupies about 21902400/1024/1024=21GB of disk space.
7. References
- Design Ideas for Database and Table Sharding in a Billion-Order System!
- Dianping Order System Database and Table Sharding Practice - Meituan Technical Team
- docs/high-concurrency/database-shard-method.md · shishan100/Java-Interview-Advanced - Gitee.com
- A Simple and Practical Database and Table Sharding Scheme Supporting Dynamic Scale-Out and Scale-In (internal link redacted)
- Interviewer: For the B+ Tree Corresponding to the Primary-Key Index of the MySQL InnoDB Storage Engine, Assuming a Page Size of 64KB, an int Primary Key, and 1KB per Row, How Many Records Can Three Levels Hold?
- Interviewer: MySQL Single-Table Performance Drops Severely at 20 Million Rows. Why? I: Uh, I Don’t Know… - Zhihu
- MySQL Single-Table Data Should Not Exceed 5 Million Rows: Experience Value or Golden Rule? - SegmentFault

Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub