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.

DatabasesCreated Updated 5 min readhistorical

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

  1. Use single-table queries and wrap the results in the code layer.
  2. Use a wide table, making writes heavier and reads lighter.
  3. 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

Discussion

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