Skip to content
JackSparrow414
Go back

Using ShardingSphere–ShardingJDBC (Part 1): Data Sharding

Table of contents

Open Table of contents

Article body

Background:

In development, we often have a few large business tables: user tables, order tables, or other central application tables. Here, large means a huge volume of data. These tables may quickly grow to millions, tens of millions, or hundreds of millions of rows, and continue growing rapidly. A single table may no longer meet storage requirements, and queries can be slow even with sensible indexes. We then need database and table sharding.

For example:

Suppose the user table is growing rapidly. We can shard its database and tables, scaling out through multiple MySQL instances. Each instance has a user database, and each database contains N user tables named user_N (N=0,1,…). Their schemas are identical; only the table names differ.

Basic principles for database and table sharding:

I recommend keeping each table below 10 million rows.

Choosing the number of table shards:

N = (total rows expected over the next 3–5 years) / recommended rows per table (10 million).

N is the number of table shards in each database. I recommend allocating enough shards up front. Expanding because there are too few is troublesome, especially when the sharding key uses a hash algorithm.

Choosing the number of database shards:

Calculate from storage capacity: (storage needed over the next 3–5 years) / recommended storage per database (less than 300 GB).

ShardingSphere website

I chose ShardingSphere because many earlier database middleware projects, such as Mycat and TDDL, were no longer maintained in my experience, whereas ShardingSphere remained active and graduated to an Apache top-level project this year. Meituan’s open-source Zebra has configuration that is too complicated compared with ShardingSphere.

If database middleware design interests you, read this Meituan article.

Database middleware has two proxy approaches: a client-side proxy (DataSource) and a server-side proxy (database proxy). The following descriptions are quoted from the Meituan article:

ShardingSphere-JDBC is a client-side data-source proxy.

ShardingSphere collectively calls database and table sharding data sharding.

  1. Add the dependency.
        <dependency>
            <groupId>org.apache.shardingsphere</groupId>
            <artifactId>sharding-jdbc-spring-boot-starter</artifactId>
            <version>4.1.0</version>
        </dependency>
  1. Configure it.
server:
  port: 8999
spring:
  application:
    name: mybatis-demo
  shardingsphere:
    datasource:
      # Database aliases
      names: ds0,ds1
      ds0:
        type: com.alibaba.druid.pool.DruidDataSource
        driverClassName: com.mysql.jdbc.Driver
        url: jdbc:mysql://localhost:3306/dhb?serverTimezone=UTC
        password: 12345
        username: root
      ds1:
        type: com.alibaba.druid.pool.DruidDataSource
        driverClassName: com.mysql.jdbc.Driver
        url: jdbc:mysql://localhost:3310/dhb?serverTimezone=UTC
        password: 12345
        username: root
    sharding:
      # Default database sharding strategy
      default-database-strategy:
        inline:
          sharding-column: id
          algorithm-expression: ds$->{id % 2}
      # Default table sharding strategy
      default-table-strategy:
        inline:
          sharding-column: age
          algorithm-expression: user_$->{age % 2}
      # Data nodes
      tables:
        user:
          actual-data-nodes: ds$->{0..1}.user_$->{0..1}
      # Default database
      default-data-source-name: ds0
    props:
      # Print SQL.
      sql.show: true
      check:
        table:
          metadata:
          # Whether to check table-shard metadata consistency at startup
          enabled: true
    # Add this setting because Druid conflicts with the default data source.
  main:
    allow-bean-definition-overriding: true

Configuration explanation:

There are two local MySQL instances, on 3306 and 3310. Each has a dhb database containing two user tables, user_0 and user_1:

Navicat showing user_0 and user_1 sharded tables in two MySQL instances

How the configuration works:

The database sharding strategy uses id % 2. A result of 0 selects ds0, the database in the 3306 instance; a result of 1 selects ds1 in the 3310 instance.

After selecting the database, which table is chosen? age % 2 selects user_0 for 0, or user_1 for 1.

For logical tables, actual tables, and data nodes, see the official documentation.

The configured strategy is inline-expression sharding. In a real application, you can implement the official interfaces to provide database and table sharding algorithms for your business. See sharding for the four interfaces: StandardShardingStrategy, ComplexShardingStrategy, and HintShardingStrategy.

For inline-expression syntax, see the official examples.

Important note:

If Alibaba’s Druid data source causes a data-source conflict at startup, use the final line of the configuration above.

  1. The entity class uses Lombok and MyBatis-Plus to simplify operations. The table name here is the logical table user, corresponding to four actual tables: ds0.user_0, ds0.user_1, ds1.user_0, and ds1.user_1.
@Data
@TableName(value = "user")
public class UserItem {
    @TableId
    private Long id;
    private String name;
    private Integer age;
    private Integer del;
}
  1. Use it to insert a record.

The console prints the following SQL:

ShardingSphere SQL logs showing an INSERT routed to ds0.user_1

ShardingSphere-JDBC writes the data to user_1 in ds0, the 3306 instance, according to the configured database and table sharding strategies.

A query containing the sharding key can locate a particular table in a particular database. Without the sharding key, it performs full routing, querying all four destinations.

Based on the configured sharding strategy, ShardingSphere-JDBC performs SQL parsing => execution optimization => SQL routing => SQL rewriting => SQL execution => result merging. See the internal architecture.

This is only a small usage example. Real business scenarios are more complicated, including sharding parent-child tables and joining sharded with unsharded tables.


Share this post:

Continue this series

ShardingSphere / ShardingJDBC

  1. Using ShardingSphere–ShardingJDBC (Part 1): Data ShardingYou are here
  2. Using ShardingSphere Database Middleware: ShardingJDBC (Part 2), Read/Write Splitting
  3. Using ShardingSphere–ShardingJDBC (Part 3): Data Masking
  4. ShardingSphere (Part 4): Custom Encryption Strategies for Data Masking

Comments

Questions, corrections, and experiences are welcome. Sign in with GitHub to comment; both language versions share this discussion.

Comments are available on the live site only.