1. 引言
1.1 分库分表的背景意义
随着互联网技术的快速发展和业务规模的不断扩张,传统单库单表架构面临着前所未有的挑战:
- 数据量爆炸性增长:单表数据量突破千万级后,查询性能呈指数级下降
- 高并发访问压力:读写请求集中在单库,导致连接池耗尽、响应时间过长
- 存储空间瓶颈:单个数据库实例的存储容量有限,无法满足海量数据存储需求
- 运维扩容困难:单库架构难以进行水平扩展,垂直扩展成本高昂
分库分表作为解决上述问题的经典方案,通过将数据分散存储到多个数据库实例和表中,实现了:
- 性能提升:减少单表数据量,显著提升查询效率
- 水平扩展:支持动态增加节点,应对业务增长
- 高可用性:分散存储降低了单点故障影响范围
1.2 ShardingSphere的优势
Apache ShardingSphere是一款开源的分布式数据库中间件生态圈,由当当网开源的Sharding-JDBC发展而来,具有以下核心优势:
技术优势:
- 透明化接入:对业务代码零侵入,支持JDBC、Proxy两种接入方式
- 丰富的分片策略:支持标准分片、复合分片、Hint分片等多种策略
- 完善的生态支持:提供分布式事务、读写分离、数据加密等全套功能
- 可插拔架构:基于微内核架构,支持功能模块的灵活扩展
版本特性:
- ShardingSphere 5.x版本采用了"连接、增强、可插拔"的设计哲学
- 支持多模态异构数据库(MySQL、PostgreSQL、Oracle等)
- 新增SQL Hint强制路由、异步数据一致性校验等高级特性
- 去除Spring配置,统一使用YAML配置,简化部署流程
1.3 本文的目标价值
本文将从实战角度出发,通过完整的电商订单系统示例,深入讲解:
- 如何根据业务场景设计合理的分库分表方案
- ShardingSphere 5.x版本的完整配置和代码实现
- 分布式事务的处理方案和最佳实践
- 性能优化策略和常见问题解决方案
通过本文的学习,读者能够掌握在企业级项目中应用ShardingSphere的核心技能,为应对海量数据存储和高并发访问提供可靠的技术保障。
2. 核心概念
2.1 数据分片
数据分片是ShardingSphere的核心功能,分为水平分片和垂直分片:
水平分片:
- 分库:将数据分散到不同的数据库实例中
- 分表:将同一数据库的数据分散到不同的物理表中
- 分片键:用于决定数据路由策略的字段,如用户ID、订单ID等
- 分片算法:基于分片键计算目标数据源和表的算法
垂直分片:
- 将表按照业务功能拆分到不同的数据库
- 如用户库、订单库、商品库等
分片策略对比:
| 标准分片 | 单一分片键 | 配置简单,路由高效 | 不支持复杂查询 |
| 复合分片 | 多个分片键 | 支持复杂查询条件 | 配置相对复杂 |
| Hint分片 | 分片键不在SQL中 | 灵活控制路由 | 需要业务代码配合 |
2.2 读写分离
读写分离通过主从复制架构实现,核心机制:
- 写操作:路由到主库(Master)
- 读操作:路由到从库(Slave),支持负载均衡
- 事务一致性:事务内的读操作默认路由到主库,避免脏读
负载均衡策略:
- 轮询(ROUND_ROBIN):按顺序依次访问从库
- 随机(RANDOM):随机选择从库
- 自定义策略:根据业务需求实现个性化负载均衡
2.3 分布式事务
在分库分表环境下,跨库操作需要分布式事务保证数据一致性,ShardingSphere支持:
XA事务(强一致性):
- 基于两阶段提交协议(2PC)
- 适用于对一致性要求极高的金融场景
- 性能相对较低,存在全局锁问题
SEATA事务(最终一致性):
- 基于Seata的AT模式,通过全局锁和回滚日志实现
- 性能优于XA,适合大多数业务场景
- 需要部署Seata Server,增加运维复杂度
本地事务:
- 在单一分片内的操作使用本地事务
- 性能最优,但无法保证跨库操作的原子性
2.4 其他核心概念
逻辑表与物理表:
- 逻辑表:应用程序操作的逻辑表名,如t_order
- 物理表:实际存储数据的表,如t_order_0、t_order_1
广播表:
- 在每个分片数据源中都存在的表
- 数据完全一致,适用于字典表、配置表等小表
绑定表:
- 分片规则一致的主表和子表
- 能够避免跨库JOIN,提升查询性能
3. 环境准备
3.1 开发环境配置
基础环境:
JDK版本: JDK 1.8+
构建工具: Maven 3.6+
数据库: MySQL 8.0+
应用框架: Spring Boot 2.7.6
ORM框架: MyBatis–Plus 3.5.3
ShardingSphere版本选择:
推荐版本: ShardingSphere 5.5.2
下载地址: https://shardingsphere.apache.org/document/5.5.2/en/downloads/
3.2 Maven依赖配置
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.7.6</version>
</parent>
<groupId>com.example</groupId>
<artifactId>shardingsphere-demo</artifactId>
<version>1.0.0</version>
<properties>
<java.version>11</java.version>
<shardingsphere.version>5.5.2</shardingsphere.version>
<mybatis-plus.version>3.5.3</mybatis-plus.version>
</properties>
<dependencies>
<!– Spring Boot Starter –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!– ShardingSphere JDBC Starter –>
<dependency>
<groupId>org.apache.shardingsphere</groupId>
<artifactId>shardingsphere-jdbc-core-spring-boot-starter</artifactId>
<version>${shardingsphere.version}</version>
</dependency>
<!– MyBatis-Plus –>
<dependency>
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-boot-starter</artifactId>
<version>${mybatis-plus.version}</version>
</dependency>
<!– MySQL驱动 –>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>8.0.33</version>
</dependency>
<!– 分布式事务支持(可选) –>
<dependency>
<groupId>org.apache.shardingsphere</groupId>
<artifactId>shardingsphere-transaction-xa-core</artifactId>
<version>${shardingsphere.version}</version>
</dependency>
<!– Seata分布式事务(可选) –>
<dependency>
<groupId>io.seata</groupId>
<artifactId>seata-all</artifactId>
<version>1.7.0</version>
</dependency>
<!– Lombok –>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<!– 测试依赖 –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
3.3 数据库准备
创建分库:
— 创建订单数据库ds0
CREATE DATABASE IF NOT EXISTS `ds0`
DEFAULT CHARACTER SET utf8mb4
COLLATE utf8mb4_general_ci;
— 创建订单数据库ds1
CREATE DATABASE IF NOT EXISTS `ds1`
DEFAULT CHARACTER SET utf8mb4
COLLATE utf8mb4_general_ci;
创建订单表结构:
— 在ds0数据库中创建订单表
USE ds0;
— 订单表0
CREATE TABLE IF NOT EXISTS `t_order_0` (
`order_id` BIGINT(20) NOT NULL COMMENT '订单ID',
`user_id` BIGINT(20) NOT NULL COMMENT '用户ID',
`order_no` VARCHAR(64) NOT NULL COMMENT '订单编号',
`total_amount` DECIMAL(10,2) NOT NULL COMMENT '订单总金额',
`status` TINYINT(1) NOT NULL DEFAULT '0' COMMENT '订单状态 0-待支付 1-已支付 2-已发货 3-已完成',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`order_id`),
KEY `idx_user_id` (`user_id`),
KEY `idx_order_no` (`order_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表0';
— 订单表1
CREATE TABLE IF NOT EXISTS `t_order_1` (
`order_id` BIGINT(20) NOT NULL COMMENT '订单ID',
`user_id` BIGINT(20) NOT NULL COMMENT '用户ID',
`order_no` VARCHAR(64) NOT NULL COMMENT '订单编号',
`total_amount` DECIMAL(10,2) NOT NULL COMMENT '订单总金额',
`status` TINYINT(1) NOT NULL DEFAULT '0' COMMENT '订单状态 0-待支付 1-已支付 2-已发货 3-已完成',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`order_id`),
KEY `idx_user_id` (`user_id`),
KEY `idx_order_no` (`order_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表1';
— 在ds1数据库中创建相同的表结构
USE ds1;
— 订单表0
CREATE TABLE IF NOT EXISTS `t_order_0` (
`order_id` BIGINT(20) NOT NULL COMMENT '订单ID',
`user_id` BIGINT(20) NOT NULL COMMENT '用户ID',
`order_no` VARCHAR(64) NOT NULL COMMENT '订单编号',
`total_amount` DECIMAL(10,2) NOT NULL COMMENT '订单总金额',
`status` TINYINT(1) NOT NULL DEFAULT '0' COMMENT '订单状态',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`order_id`),
KEY `idx_user_id` (`user_id`),
KEY `idx_order_no` (`order_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表0';
— 订单表1
CREATE TABLE IF NOT EXISTS `t_order_1` (
`order_id` BIGINT(20) NOT NULL COMMENT '订单ID',
`user_id` BIGINT(20) NOT NULL COMMENT '用户ID',
`order_no` VARCHAR(64) NOT NULL COMMENT '订单编号',
`total_amount` DECIMAL(10,2) NOT NULL COMMENT '订单总金额',
`status` TINYINT(1) NOT NULL DEFAULT '0' COMMENT '订单状态',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`order_id`),
KEY `idx_user_id` (`user_id`),
KEY `idx_order_no` (`order_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表1';
创建订单明细表:
— 在ds0和ds1数据库中分别创建订单明细表
USE ds0;
CREATE TABLE IF NOT EXISTS `t_order_item_0` (
`item_id` BIGINT(20) NOT NULL COMMENT '订单明细ID',
`order_id` BIGINT(20) NOT NULL COMMENT '订单ID',
`product_id` BIGINT(20) NOT NULL COMMENT '商品ID',
`product_name` VARCHAR(128) NOT NULL COMMENT '商品名称',
`quantity` INT(11) NOT NULL COMMENT '购买数量',
`price` DECIMAL(10,2) NOT NULL COMMENT '单价',
`total_price` DECIMAL(10,2) NOT NULL COMMENT '小计',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
PRIMARY KEY (`item_id`),
KEY `idx_order_id` (`order_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单明细表0';
CREATE TABLE IF NOT EXISTS `t_order_item_1` (
`item_id` BIGINT(20) NOT NULL COMMENT '订单明细ID',
`order_id` BIGINT(20) NOT NULL COMMENT '订单ID',
`product_id` BIGINT(20) NOT NULL COMMENT '商品ID',
`product_name` VARCHAR(128) NOT NULL COMMENT '商品名称',
`quantity` INT(11) NOT NULL COMMENT '购买数量',
`price` DECIMAL(10,2) NOT NULL COMMENT '单价',
`total_price` DECIMAL(10,2) NOT NULL COMMENT '小计',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
PRIMARY KEY (`item_id`),
KEY `idx_order_id` (`order_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单明细表1';
— ds1数据库中的订单明细表结构相同,省略…
3.4 分布式事务表准备
如果使用SEATA事务,需要创建undo_log表:
— 在每个数据库中创建undo_log表
USE ds0;
CREATE TABLE IF NOT EXISTS `undo_log` (
`id` BIGINT(20) NOT NULL AUTO_INCREMENT COMMENT 'increment id',
`branch_id` BIGINT(20) NOT NULL COMMENT 'branch transaction id',
`xid` VARCHAR(100) NOT NULL COMMENT 'global transaction id',
`context` VARCHAR(128) NOT NULL COMMENT 'undo_log context',
`rollback_info` LONGBLOB NOT NULL COMMENT 'rollback info',
`log_status` INT(11) NOT NULL COMMENT '0:normal status,1:defense status',
`log_created` DATETIME NOT NULL COMMENT 'create datetime',
`log_modified` DATETIME NOT NULL COMMENT 'modify datetime',
PRIMARY KEY (`id`),
UNIQUE KEY `ux_undo_log` (`xid`, `branch_id`)
) ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8mb4 COMMENT='AT transaction mode undo table';
— ds1数据库中也创建相同的undo_log表
USE ds1;
CREATE TABLE IF NOT EXISTS `undo_log` (
— 省略,结构与ds0相同
) ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8mb4;
4. 项目实战
4.1 分库分表方案设计
4.1.1 分片策略选择
基于电商订单系统的业务特点,我们设计如下分片策略:
业务分析:
- 订单数据量:预估每年1亿条订单
- 查询模式:90%的查询基于用户ID,10%基于订单ID
- 扩展需求:支持按时间归档历史数据
分片方案:
分库策略:2个数据库(ds0, ds1)
分库键:user_id
分库算法:user_id % 2
分表策略:每个库2张表(t_order_0, t_order_1)
分表键:order_id
分表算法:order_id % 2
数据节点规划:
逻辑表:t_order
物理表分布:
– ds0.t_order_0
– ds0.t_order_1
– ds1.t_order_0
– ds1.t_order_1
实际数据节点:ds$->{0..1}.t_order_$->{0..1}
绑定表配置:
绑定表:t_order, t_order_item
目的:确保同一订单的订单和明细在同一个物理表中,避免跨库JOIN
4.1.2 分片键设计原则
分片键选择原则:
常见分片键选择场景:
- 用户表:user_id(用户维度查询最多)
- 订单表:order_id或user_id(根据查询模式选择)
- 交易记录:account_id(账户维度查询)
- 日志表:create_time(时间维度查询)
4.1.3 架构设计图
┌─────────────────────────────────────────────────────────────┐
│ 应用层 │
│ Spring Boot + MyBatis-Plus │
└──────────────────────┬──────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ ShardingSphere-JDBC │
│ (SQL解析、路由、改写、执行、结果归并) │
└──────────────────────┬──────────────────────────────────────┘
│
┌──────────────┼──────────────┐
▼ ▼ ▼
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ ds0库 │ │ ds1库 │ │ 配置中心 │
│ t_order_0 │ │ t_order_0 │ │ (可选) │
│ t_order_1 │ │ t_order_1 │ │ │
│ t_order_…│ │ t_order_…│ │ │
└─────────────┘ └─────────────┘ └─────────────┘
4.2 完整的代码实现
4.2.1 项目结构
shardingsphere-demo/
├── src/
│ ├── main/
│ │ ├── java/
│ │ │ └── com/
│ │ │ └── example/
│ │ │ └── shardingsphere/
│ │ │ ├── ShardingSphereDemoApplication.java
│ │ │ ├── config/
│ │ │ │ └── MybatisPlusConfig.java
│ │ │ ├── entity/
│ │ │ │ ├── Order.java
│ │ │ │ └── OrderItem.java
│ │ │ ├── mapper/
│ │ │ │ ├── OrderMapper.java
│ │ │ │ └── OrderItemMapper.java
│ │ │ ├── service/
│ │ │ │ ├── OrderService.java
│ │ │ │ └── impl/
│ │ │ │ └── OrderServiceImpl.java
│ │ │ └── controller/
│ │ │ └── OrderController.java
│ │ └── resources/
│ │ ├── application.yml
│ │ └── logback-spring.xml
│ └── test/
│ └── java/
│ └── com/
│ └── example/
│ └── shardingsphere/
│ └── ShardingSphereTest.java
└── pom.xml
4.2.2 核心配置文件
application.yml配置:
# ==================== 服务基础配置 ====================
server:
port: 8080 # 应用服务监听的端口号
servlet:
context-path: /api # 应用访问的上下文路径,完整的访问路径为 http://localhost:8080/api
# ==================== Spring应用配置 ====================
spring:
application:
name: shardingsphere–demo # 应用名称,用于标识和监控
# ==================== ShardingSphere分库分表核心配置 ====================
shardingsphere:
# ==================== 数据源配置区域 ====================
# 说明:配置物理数据库连接信息,ShardingSphere会基于这些数据源进行分片路由
datasource:
# 数据源名称列表,用逗号分隔,这些名称在后续配置中会被引用
names: ds0,ds1
# ==================== ds0数据源配置(第一个数据库) ====================
ds0:
# 数据源类型,使用HikariCP高性能连接池
type: com.zaxxer.hikari.HikariDataSource
# MySQL驱动类名,注意JDK1.8环境使用MySQL 8.0+驱动时使用com.mysql.cj.jdbc.Driver
driver-class-name: com.mysql.cj.jdbc.Driver
# 数据库连接URL,参数说明:
# – localhost:3306 数据库服务器地址和端口
# – ds0 数据库名称(物理数据库)
# – useSSL=false 禁用SSL连接,本地开发环境常用
# – serverTimezone=UTC 设置服务器时区为UTC,避免时区转换问题
# – useUnicode=true 启用Unicode字符集支持
# – characterEncoding=utf8 设置字符编码为UTF-8,支持中文
# – allowPublicKeyRetrieval=true 允许公钥检索,解决某些MySQL版本的连接问题
jdbc-url: jdbc:mysql://localhost:3306/ds0?useSSL=false&serverTimezone=UTC&useUnicode=true&characterEncoding=utf8&allowPublicKeyRetrieval=true
username: root # 数据库用户名
password: root # 数据库密码
# ==================== HikariCP连接池详细配置 ====================
hikari:
minimum-idle: 5 # 连接池中最小空闲连接数,即使没有请求也会保持这些连接,避免频繁创建销毁
maximum-pool-size: 20 # 连接池中最大连接数,根据数据库性能和并发量调整,一般设置为CPU核心数*2
connection-timeout: 30000 # 连接超时时间(毫秒),30秒内无法获取连接会抛出异常
idle-timeout: 600000 # 空闲连接超时时间(毫秒),10分钟未使用的连接会被回收,避免资源浪费
max-lifetime: 1800000 # 连接最大生命周期(毫秒),30分钟后连接会被强制回收并重新创建,防止长时间使用导致的内存泄漏
connection-test-query: SELECT 1 # 连接测试SQL,在获取连接时执行,确保连接可用性
pool-name: ShardingHikariCP–ds0 # 连接池名称,便于监控和日志识别
# ==================== ds1数据源配置(第二个数据库) ====================
ds1:
# 数据源类型,同样使用HikariCP连接池
type: com.zaxxer.hikari.HikariDataSource
# MySQL驱动类名,必须与ds0保持一致
driver-class-name: com.mysql.cj.jdbc.Driver
# 数据库连接URL,连接到ds1数据库,其他参数与ds0相同
jdbc-url: jdbc:mysql://localhost:3306/ds1?useSSL=false&serverTimezone=UTC&useUnicode=true&characterEncoding=utf8&allowPublicKeyRetrieval=true
username: root # 数据库用户名
password: root # 数据库密码
# HikariCP连接池配置,参数与ds0保持一致
hikari:
minimum-idle: 5
maximum-pool-size: 20
connection-timeout: 30000
idle-timeout: 600000
max-lifetime: 1800000
connection-test-query: SELECT 1
pool-name: ShardingHikariCP–ds1
# ==================== 分片规则配置区域 ====================
# 说明:配置如何将逻辑表路由到物理表的核心规则
rules:
sharding:
# ==================== 分片表配置 ====================
# 说明:定义每个逻辑表对应的物理表分布和分片策略
tables:
# ==================== t_order订单表配置 ====================
t_order:
# 实际数据节点配置
# 说明:定义逻辑表t_order对应的物理表分布
# 格式:数据源名称.表名 的Groovy行表达式
# ds$->{0..1} 表示数据源ds0和ds1
# t_order_$->{0..1} 表示表名t_order_0和t_order_1
# 展开后的结果:ds0.t_order_0, ds0.t_order_1, ds1.t_order_0, ds1.t_order_1
actual-data-nodes: ds$–>{0..1}.t_order_$–>{0..1}
# ==================== 数据库分片策略配置 ====================
# 说明:配置如何根据分片键将数据路由到不同的数据库
database-strategy:
# 标准分片策略,适用于单个分片键的场景
standard:
# 分库分片键字段名,使用user_id作为分库依据
# 说明:用户ID作为分库键,确保同一用户的订单在同一个数据库中,便于用户维度的查询
sharding-column: user_id
# 分库算法名称,对应下面sharding-algorithms中的db-mod算法
sharding-algorithm-name: db–mod
# ==================== 表分片策略配置 ====================
# 说明:配置如何根据分片键将数据路由到同一数据库中的不同表
table-strategy:
# 标准分片策略
standard:
# 分表分片键字段名,使用order_id作为分表依据
# 说明:订单ID作为分表键,确保数据在同一库的不同表中均匀分布
sharding-column: order_id
# 分表算法名称,对应下面sharding-algorithms中的table-mod算法
sharding-algorithm-name: table–mod
# ==================== 主键生成策略配置 ====================
# 说明:配置分布式主键生成器,避免分库分表后主键冲突
key-generate-strategy:
# 主键字段名
column: order_id
# 主键生成器名称,对应下面key-generators中的snowflake算法
key-generator-name: snowflake
# ==================== t_order_item订单明细表配置 ====================
t_order_item:
# 实际数据节点配置
actual-data-nodes: ds$–>{0..1}.t_order_item_$–>{0..1}
# 数据库分片策略
database-strategy:
standard:
# 使用order_id作为分库键
# 说明:订单明细表与订单表使用相同的分库键,确保同一订单的明细数据与订单在同一数据库中
sharding-column: order_id
# 使用专门的分库算法db-mod-by-order
sharding-algorithm-name: db–mod–by–order
# 表分片策略
table-strategy:
standard:
# 同样使用order_id作为分表键
sharding-column: order_id
sharding-algorithm-name: table–mod
# 主键生成策略
key-generate-strategy:
column: item_id # 订单明细表的主键字段
key-generator-name: snowflake
# ==================== 绑定表配置 ====================
# 说明:绑定表是指分片规则完全一致的主表和子表
# 作用:绑定表在查询时会进行优化,避免笛卡尔积,提升JOIN查询性能
# 原理:ShardingSphere知道这些表的分片规则相同,会根据分片键将路由到同一物理表,避免跨库JOIN
binding-tables:
– t_order,t_order_item # 订单表和订单明细表为绑定表,它们的分片规则一致
# ==================== 广播表配置 ====================
# 说明:广播表是指在所有分片数据源中都存在的表,每个数据源的数据完全相同
# 适用场景:数据量小、变动不频繁的字典表、配置表、地区表等
# 优势:查询时只从一个数据源读取即可,避免跨分片查询
broadcast-tables:
– t_dict # 字典表作为广播表,在ds0和ds1中都存在且数据一致
# ==================== 分片算法配置 ====================
# 说明:定义各种分片算法的具体实现和参数
sharding-algorithms:
# ==================== 数据库分片算法-按用户ID取模 ====================
db-mod:
# 算法类型:MOD取模算法
# 说明:基于取模运算的分片算法,简单高效,适用于均匀分布数据
type: MOD
# 算法属性配置
props:
# 分片数量,表示将数据均匀分片到2个数据库
# 计算公式:user_id % 2 = 目标数据库索引(0或1)
sharding-count: 2
# ==================== 表分片算法-按订单ID取模 ====================
table-mod:
# 算法类型:MOD取模算法
type: MOD
props:
# 分片数量,表示将数据均匀分片到2张表
# 计算公式:order_id % 2 = 目标表索引(0或1)
sharding-count: 2
# ==================== 数据库分片算法-按订单ID取模(用于订单明细) ====================
db-mod-by-order:
# 算法类型:MOD取模算法
type: MOD
props:
# 分片数量为2,与订单表的分库策略保持一致
# 计算公式:order_id % 2 = 目标数据库索引(0或1)
# 说明:确保订单明细表与订单表分片到相同的数据库
sharding-count: 2
# ==================== 主键生成器配置 ====================
# 说明:定义分布式主键生成器的具体实现
key-generators:
snowflake:
# 算法类型:SNOWFLAKE雪花算法
# 说明:Twitter开源的分布式ID生成算法,能够生成全局唯一、趋势递增的64位整数ID
# 优势:不依赖数据库,性能高,支持分布式环境
type: SNOWFLAKE
# 雪花算法属性配置
props:
# 工作机器ID,取值范围0-1023
# 说明:在分布式环境中,每个节点需要设置不同的worker-id,避免ID冲突
# 单机环境可以随意设置,集群环境需要统一管理
worker-id: 123
# ==================== ShardingSphere属性配置 ====================
# 说明:配置ShardingSphere框架的全局属性和功能开关
props:
# 显示SQL语句开关
# true表示在日志中打印出路由后的实际SQL,方便调试和问题排查
# false表示不打印SQL,生产环境建议关闭以提升性能
sql-show: true
# 启用元数据检查
# true表示应用启动时检查物理表是否存在,表结构是否与配置一致
# false表示不进行元数据检查,提升启动速度但可能隐藏配置错误
# 建议:开发环境开启,生产环境可根据需要选择
check-table-metadata-enabled: true
# SQL打印日志级别(可选配置)
# 控制SQL日志的详细程度,可选值:INFO, WARN, ERROR
# sql-log-level: INFO
# 查询结果缓存大小(可选配置)
# 设置查询结果缓存的最大条数,可以提升重复查询的性能
# query-result-cache-max-size: 100
# 执行引擎类型(可选配置)
# SEMI_TRANSACTIONAL:半事务型执行引擎,适用于OLTP场景
# MEMORY_STRICTLY:内存严格型执行引擎,适用于OLAP场景
# execute-engine-type: SEMI_TRANSACTIONAL
# ==================== MyBatis-Plus配置 ====================
mybatis-plus:
# ==================== MyBatis配置 ====================
configuration:
# 开启驼峰命名转换
# 说明:将数据库字段名user_name自动转换为Java属性名userName
map-underscore-to-camel-case: true
# SQL日志实现类
# 说明:使用标准输出打印SQL语句,便于开发调试,生产环境建议改为slf4j
log-impl: org.apache.ibatis.logging.stdout.StdOutImpl
# 开启二级缓存(可选)
# cache-enabled: true
# 开启懒加载(可选)
# lazy-loading-enabled: true
# ==================== 全局配置 ====================
global-config:
# 数据库配置
db-config:
# 逻辑删除字段名
# 说明:删除操作不真正删除数据,而是将deleted字段设置为1
logic-delete-field: deleted
# 逻辑删除值:表示数据已删除
logic-delete-value: 1
# 逻辑未删除值:表示数据正常
logic-not-delete-value: 0
# 主键类型:数据库自增
# 注意:由于使用了ShardingSphere的雪花算法,这里可以设置为ASSIGN_ID
id-type: ASSIGN_ID
# 表名前缀(可选)
# table-prefix: t_
# ==================== 日志配置 ====================
logging:
level:
# 应用日志级别
# debug级别可以查看详细的SQL路由信息
com.example.shardingsphere: debug
# ShardingSphere框架日志级别
# info级别可以看到SQL路由和执行的基本信息
org.apache.shardingsphere: info
# HikariCP连接池日志级别(可选)
# com.zaxxer.hikari: debug
4.2.3 实体类
Order实体类:
package com.example.shardingsphere.entity;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.math.BigDecimal;
import java.util.Date;
/**
* 订单实体类
*/
@Data
@TableName("t_order")
public class Order {
/**
* 订单ID(主键)
*/
private Long orderId;
/**
* 用户ID(分库键)
*/
private Long userId;
/**
* 订单编号
*/
private String orderNo;
/**
* 订单总金额
*/
private BigDecimal totalAmount;
/**
* 订单状态 0-待支付 1-已支付 2-已发货 3-已完成
*/
private Integer status;
/**
* 创建时间
*/
private Date createTime;
/**
* 更新时间
*/
private Date updateTime;
}
OrderItem实体类:
package com.example.shardingsphere.entity;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.math.BigDecimal;
import java.util.Date;
/**
* 订单明细实体类
*/
@Data
@TableName("t_order_item")
public class OrderItem {
/**
* 订单明细ID(主键)
*/
private Long itemId;
/**
* 订单ID(分片键)
*/
private Long orderId;
/**
* 商品ID
*/
private Long productId;
/**
* 商品名称
*/
private String productName;
/**
* 购买数量
*/
private Integer quantity;
/**
* 单价
*/
private BigDecimal price;
/**
* 小计
*/
private BigDecimal totalPrice;
/**
* 创建时间
*/
private Date createTime;
}
4.2.4 DAO层
OrderMapper接口:
package com.example.shardingsphere.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.example.shardingsphere.entity.Order;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import org.apache.ibatis.annotations.Select;
import java.util.List;
import java.util.Map;
/**
* 订单Mapper接口
*/
@Mapper
public interface OrderMapper extends BaseMapper<Order> {
/**
* 根据订单编号查询订单
*/
@Select("SELECT * FROM t_order WHERE order_no = #{orderNo}")
Order selectByOrderNo(@Param("orderNo") String orderNo);
/**
* 统计用户订单数量
*/
@Select("SELECT COUNT(*) FROM t_order WHERE user_id = #{userId}")
int countByUserId(@Param("userId") Long userId);
/**
* 查询用户订单列表
*/
@Select("SELECT * FROM t_order WHERE user_id = #{userId} ORDER BY create_time DESC LIMIT #{limit}")
List<Order> selectUserOrders(@Param("userId") Long userId, @Param("limit") int limit);
/**
* 统计各分片订单数量
*/
@Select("SELECT SUBSTRING_INDEX(SUBSTRING_INDEX(TABLE_NAME, '_', -1), '.', 1) as shard, COUNT(*) as count " +
"FROM information_schema.tables WHERE table_schema = DATABASE() AND table_name LIKE 't_order_%' " +
"GROUP BY shard")
List<Map<String, Object>> countOrdersByShard();
}
OrderItemMapper接口:
package com.example.shardingsphere.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.example.shardingsphere.entity.OrderItem;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import org.apache.ibatis.annotations.Select;
import java.util.List;
/**
* 订单明细Mapper接口
*/
@Mapper
public interface OrderItemMapper extends BaseMapper<OrderItem> {
/**
* 根据订单ID查询订单明细
*/
@Select("SELECT * FROM t_order_item WHERE order_id = #{orderId}")
List<OrderItem> selectByOrderId(@Param("orderId") Long orderId);
/**
* 统计订单商品数量
*/
@Select("SELECT SUM(quantity) FROM t_order_item WHERE order_id = #{orderId}")
Integer countProductsByOrderId(@Param("orderId") Long orderId);
}
4.2.5 Service层
OrderService接口:
package com.example.shardingsphere.service;
import com.baomidou.mybatisplus.extension.service.IService;
import com.example.shardingsphere.entity.Order;
import com.example.shardingsphere.entity.OrderItem;
import java.util.List;
import java.util.Map;
/**
* 订单服务接口
*/
public interface OrderService extends IService<Order> {
/**
* 创建订单(含订单明细)
*/
boolean createOrder(Order order, List<OrderItem> orderItems);
/**
* 根据订单ID查询订单及明细
*/
Map<String, Object> getOrderWithItems(Long orderId);
/**
* 根据用户ID查询订单列表
*/
List<Order> getOrdersByUserId(Long userId);
/**
* 更新订单状态
*/
boolean updateOrderStatus(Long orderId, Integer status);
/**
* 删除订单(逻辑删除)
*/
boolean deleteOrder(Long orderId);
/**
* 统计分片数据分布
*/
Map<String, Object> getShardingStatistics();
}
OrderServiceImpl实现类:
package com.example.shardingsphere.service.impl;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.example.shardingsphere.entity.Order;
import com.example.shardingsphere.entity.OrderItem;
import com.example.shardingsphere.mapper.OrderItemMapper;
import com.example.shardingsphere.mapper.OrderMapper;
import com.example.shardingsphere.service.OrderService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
* 订单服务实现类
*/
@Slf4j
@Service
public class OrderServiceImpl extends ServiceImpl<OrderMapper, Order> implements OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private OrderItemMapper orderItemMapper;
/**
* 创建订单(含订单明细)
*/
@Override
@Transactional(rollbackFor = Exception.class)
public boolean createOrder(Order order, List<OrderItem> orderItems) {
try {
log.info("开始创建订单,用户ID:{},订单编号:{}", order.getUserId(), order.getOrderNo());
// 1. 保存订单
int orderResult = orderMapper.insert(order);
if (orderResult <= 0) {
log.error("订单插入失败,订单编号:{}", order.getOrderNo());
throw new RuntimeException("订单插入失败");
}
log.info("订单插入成功,订单ID:{}", order.getOrderId());
// 2. 保存订单明细
for (OrderItem orderItem : orderItems) {
orderItem.setOrderId(order.getOrderId());
// 计算小计
orderItem.setTotalPrice(orderItem.getPrice().multiply(
new java.math.BigDecimal(orderItem.getQuantity())
));
int itemResult = orderItemMapper.insert(orderItem);
if (itemResult <= 0) {
log.error("订单明细插入失败,订单ID:{}", order.getOrderId());
throw new RuntimeException("订单明细插入失败");
}
}
log.info("订单创建成功,订单ID:{},明细数量:{}", order.getOrderId(), orderItems.size());
return true;
} catch (Exception e) {
log.error("创建订单失败,订单编号:{},错误信息:{}", order.getOrderNo(), e.getMessage(), e);
throw new RuntimeException("创建订单失败:" + e.getMessage());
}
}
/**
* 根据订单ID查询订单及明细
*/
@Override
public Map<String, Object> getOrderWithItems(Long orderId) {
log.info("查询订单及明细,订单ID:{}", orderId);
Map<String, Object> result = new HashMap<>();
// 查询订单信息
Order order = orderMapper.selectById(orderId);
if (order == null) {
log.warn("订单不存在,订单ID:{}", orderId);
return null;
}
result.put("order", order);
// 查询订单明细
List<OrderItem> orderItems = orderItemMapper.selectByOrderId(orderId);
result.put("orderItems", orderItems);
log.info("查询订单及明细成功,订单ID:{},明细数量:{}", orderId, orderItems.size());
return result;
}
/**
* 根据用户ID查询订单列表
*/
@Override
public List<Order> getOrdersByUserId(Long userId) {
log.info("查询用户订单列表,用户ID:{}", userId);
List<Order> orders = orderMapper.selectUserOrders(userId, 100);
log.info("查询用户订单列表成功,用户ID:{},订单数量:{}", userId, orders.size());
return orders;
}
/**
* 更新订单状态
*/
@Override
@Transactional(rollbackFor = Exception.class)
public boolean updateOrderStatus(Long orderId, Integer status) {
log.info("更新订单状态,订单ID:{},状态:{}", orderId, status);
Order order = new Order();
order.setOrderId(orderId);
order.setStatus(status);
int result = orderMapper.updateById(order);
if (result > 0) {
log.info("更新订单状态成功,订单ID:{}", orderId);
return true;
} else {
log.warn("更新订单状态失败,订单ID:{}", orderId);
return false;
}
}
/**
* 删除订单(逻辑删除)
*/
@Override
@Transactional(rollbackFor = Exception.class)
public boolean deleteOrder(Long orderId) {
log.info("删除订单,订单ID:{}", orderId);
int result = orderMapper.deleteById(orderId);
if (result > 0) {
log.info("删除订单成功,订单ID:{}", orderId);
return true;
} else {
log.warn("删除订单失败,订单ID:{}", orderId);
return false;
}
}
/**
* 统计分片数据分布
*/
@Override
public Map<String, Object> getShardingStatistics() {
log.info("统计分片数据分布");
Map<String, Object> statistics = new HashMap<>();
// 统计各分片订单数量
List<Map<String, Object>> shardCounts = orderMapper.countOrdersByShard();
statistics.put("shardCounts", shardCounts);
// 统计总订单数
int totalOrders = orderMapper.selectCount(null);
statistics.put("totalOrders", totalOrders);
log.info("统计分片数据分布成功,总订单数:{}", totalOrders);
return statistics;
}
}
4.2.6 Controller层
OrderController控制器:
package com.example.shardingsphere.controller;
import com.example.shardingsphere.entity.Order;
import com.example.shardingsphere.entity.OrderItem;
import com.example.shardingsphere.service.OrderService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
* 订单控制器
*/
@Slf4j
@RestController
@RequestMapping("/orders")
public class OrderController {
@Autowired
private OrderService orderService;
/**
* 创建订单
*/
@PostMapping("/create")
public Map<String, Object> createOrder(@RequestBody Map<String, Object> request) {
Map<String, Object> response = new HashMap<>();
try {
// 解析请求参数
Order order = parseOrder(request);
List<OrderItem> orderItems = parseOrderItems(request);
// 创建订单
boolean success = orderService.createOrder(order, orderItems);
if (success) {
response.put("success", true);
response.put("message", "订单创建成功");
response.put("orderId", order.getOrderId());
} else {
response.put("success", false);
response.put("message", "订单创建失败");
}
} catch (Exception e) {
log.error("创建订单异常", e);
response.put("success", false);
response.put("message", "创建订单异常:" + e.getMessage());
}
return response;
}
/**
* 查询订单详情
*/
@GetMapping("/{orderId}")
public Map<String, Object> getOrder(@PathVariable Long orderId) {
Map<String, Object> response = new HashMap<>();
try {
Map<String, Object> orderData = orderService.getOrderWithItems(orderId);
if (orderData != null) {
response.put("success", true);
response.put("data", orderData);
} else {
response.put("success", false);
response.put("message", "订单不存在");
}
} catch (Exception e) {
log.error("查询订单异常", e);
response.put("success", false);
response.put("message", "查询订单异常:" + e.getMessage());
}
return response;
}
/**
* 查询用户订单列表
*/
@GetMapping("/user/{userId}")
public Map<String, Object> getUserOrders(@PathVariable Long userId) {
Map<String, Object> response = new HashMap<>();
try {
List<Order> orders = orderService.getOrdersByUserId(userId);
response.put("success", true);
response.put("data", orders);
response.put("total", orders.size());
} catch (Exception e) {
log.error("查询用户订单异常", e);
response.put("success", false);
response.put("message", "查询用户订单异常:" + e.getMessage());
}
return response;
}
/**
* 更新订单状态
*/
@PutMapping("/{orderId}/status")
public Map<String, Object> updateOrderStatus(
@PathVariable Long orderId,
@RequestParam Integer status) {
Map<String, Object> response = new HashMap<>();
try {
boolean success = orderService.updateOrderStatus(orderId, status);
if (success) {
response.put("success", true);
response.put("message", "订单状态更新成功");
} else {
response.put("success", false);
response.put("message", "订单状态更新失败");
}
} catch (Exception e) {
log.error("更新订单状态异常", e);
response.put("success", false);
response.put("message", "更新订单状态异常:" + e.getMessage());
}
return response;
}
/**
* 删除订单
*/
@DeleteMapping("/{orderId}")
public Map<String, Object> deleteOrder(@PathVariable Long orderId) {
Map<String, Object> response = new HashMap<>();
try {
boolean success = orderService.deleteOrder(orderId);
if (success) {
response.put("success", true);
response.put("message", "订单删除成功");
} else {
response.put("success", false);
response.put("message", "订单删除失败");
}
} catch (Exception e) {
log.error("删除订单异常", e);
response.put("success", false);
response.put("message", "删除订单异常:" + e.getMessage());
}
return response;
}
/**
* 统计分片数据分布
*/
@GetMapping("/statistics")
public Map<String, Object> getStatistics() {
Map<String, Object> response = new HashMap<>();
try {
Map<String, Object> statistics = orderService.getShardingStatistics();
response.put("success", true);
response.put("data", statistics);
} catch (Exception e) {
log.error("统计数据异常", e);
response.put("success", false);
response.put("message", "统计数据异常:" + e.getMessage());
}
return response;
}
/**
* 解析订单信息
*/
private Order parseOrder(Map<String, Object> request) {
Order order = new Order();
order.setUserId(Long.valueOf(request.get("userId").toString()));
order.setOrderNo(request.get("orderNo").toString());
order.setTotalAmount(new java.math.BigDecimal(request.get("totalAmount").toString()));
order.setStatus(Integer.valueOf(request.get("status").toString()));
return order;
}
/**
* 解析订单明细
*/
@SuppressWarnings("unchecked")
private List<OrderItem> parseOrderItems(Map<String, Object> request) {
List<Map<String, Object>> itemsData =
(List<Map<String, Object>>) request.get("orderItems");
return itemsData.stream().map(itemData -> {
OrderItem item = new OrderItem();
item.setProductId(Long.valueOf(itemData.get("productId").toString()));
item.setProductName(itemData.get("productName").toString());
item.setQuantity(Integer.valueOf(itemData.get("quantity").toString()));
item.setPrice(new java.math.BigDecimal(itemData.get("price").toString()));
return item;
}).collect(java.util.stream.Collectors.toList());
}
}
4.3 关键配置详解
4.3.1 数据源配置
HikariCP连接池优化参数:
hikari:
minimum-idle: 5 # 最小空闲连接数
maximum-pool-size: 20 # 最大连接池大小
connection-timeout: 30000 # 连接超时时间(毫秒)
idle-timeout: 600000 # 空闲连接超时时间(毫秒)
max-lifetime: 1800000 # 连接最大生命周期(毫秒)
connection-test-query: SELECT 1 # 连接测试查询
pool-name: ShardingHikariCP # 连接池名称
连接池调优建议:
- CPU密集型应用:maximum-pool-size = CPU核心数 + 1
- IO密集型应用:maximum-pool-size = CPU核心数 * 2
- 数据库连接限制:考虑数据库的最大连接数限制
4.3.2 分片规则配置
数据库分片策略:
database-strategy:
standard:
sharding-column: user_id # 分库键
sharding-algorithm-name: db–mod # 分片算法名称
表分片策略:
table-strategy:
standard:
sharding-column: order_id # 分表键
sharding-algorithm-name: table–mod # 分片算法名称
分片算法配置:
sharding-algorithms:
db-mod:
type: MOD # 取模分片算法
props:
sharding-count: 2 # 分片数量
table-mod:
type: MOD # 取模分片算法
props:
sharding-count: 2 # 分片数量
内置分片算法类型:
| MOD | 取模分片 | 数据均匀分布 |
| HASH_MOD | 哈希取模分片 | 字符串类型分片键 |
| VOLUME_RANGE | 体积范围分片 | 按数值范围分片 |
| BOUNDARY_RANGE | 边界范围分片 | 按固定边界分片 |
| AUTO_INTERVAL | 自动间隔分片 | 按时间间隔分片 |
| INTERVAL | 间隔分片 | 按固定间隔分片 |
4.3.3 分布式事务配置
XA事务配置:
spring:
shardingsphere:
props:
# 事务类型:LOCAL, XA, BASE
transaction-type: XA
# XA事务管理器
xa-transaction-manager-class-name: org.apache.shardingsphere.transaction.xa.AtomikosXATransactionManager
SEATA事务配置:
spring:
shardingsphere:
props:
transaction-type: SEATA
# Seata相关配置
seata-enable-auto-data-source-proxy: true
Seata配置文件(seata.conf):
client {
application.id = shardingsphere-demo
transaction.service.group = my_test_tx_group
}
使用分布式事务:
@Service
public class OrderServiceImpl {
@Transactional
@org.apache.shardingsphere.transaction.annotation.ShardingTransactionType(
org.apache.shardingsphere.transaction.core.TransactionType.XA)
public void createOrderWithXA(Order order, List<OrderItem> orderItems) {
// 跨库事务操作
orderMapper.insert(order);
for (OrderItem item : orderItems) {
orderItemMapper.insert(item);
}
}
}
4.3.4 读写分离配置
读写分离完整配置:
spring:
shardingsphere:
# 数据源配置
datasource:
names: master,slave0,slave1
master:
type: com.zaxxer.hikari.HikariDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
jdbc-url: jdbc:mysql://master–host:3306/db_name
username: root
password: root
slave0:
type: com.zaxxer.hikari.HikariDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
jdbc-url: jdbc:mysql://slave0–host:3306/db_name
username: root
password: root
slave1:
type: com.zaxxer.hikari.HikariDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
jdbc-url: jdbc:mysql://slave1–host:3306/db_name
username: root
password: root
# 读写分离规则
rules:
readwrite-splitting:
data-sources:
readwrite_ds:
type: Static
props:
write-data-source-name: master
read-data-source-names: slave0,slave1
load-balancer-name: round_robin
load-balancers:
round_robin:
type: ROUND_ROBIN
负载均衡策略配置:
load-balancers:
round_robin:
type: ROUND_ROBIN # 轮询策略
random:
type: RANDOM # 随机策略
custom:
type: CLASS_BASED # 自定义策略
props:
strategy-class-name: com.example.CustomLoadBalanceAlgorithm
5. 测试验证
5.1 测试用例代码
ShardingSphereTest测试类:
package com.example.shardingsphere;
import com.example.shardingsphere.entity.Order;
import com.example.shardingsphere.entity.OrderItem;
import com.example.shardingsphere.mapper.OrderMapper;
import com.example.shardingsphere.mapper.OrderItemMapper;
import com.example.shardingsphere.service.OrderService;
import lombok.extern.slf4j.Slf4j;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import java.math.BigDecimal;
import java.util.*;
/**
* ShardingSphere分库分表测试类
*/
@Slf4j
@SpringBootTest
public class ShardingSphereTest {
@Autowired
private OrderService orderService;
@Autowired
private OrderMapper orderMapper;
@Autowired
private OrderItemMapper orderItemMapper;
/**
* 测试订单插入
*/
@Test
public void testInsertOrder() {
log.info("===== 开始测试订单插入 =====");
for (int i = 1; i <= 10; i++) {
// 创建订单
Order order = new Order();
order.setUserId((long) (i % 2 + 1)); // 用户ID在1和2之间切换
order.setOrderNo("ORD" + System.currentTimeMillis() + i);
order.setTotalAmount(new BigDecimal("99.99"));
order.setStatus(0);
order.setCreateTime(new Date());
order.setUpdateTime(new Date());
// 创建订单明细
List<OrderItem> orderItems = new ArrayList<>();
for (int j = 1; j <= 3; j++) {
OrderItem item = new OrderItem();
item.setProductId((long) (j * 100));
item.setProductName("商品" + j);
item.setQuantity(2);
item.setPrice(new BigDecimal("10.00"));
item.setCreateTime(new Date());
orderItems.add(item);
}
// 插入订单
boolean success = orderService.createOrder(order, orderItems);
log.info("订单插入结果:{}, 订单ID:{}, 用户ID:{}, 订单编号:{}",
success, order.getOrderId(), order.getUserId(), order.getOrderNo());
}
log.info("===== 订单插入测试完成 =====");
}
/**
* 测试订单查询
*/
@Test
public void testSelectOrder() {
log.info("===== 开始测试订单查询 =====");
// 根据订单ID查询
Long orderId = 1234567890123456789L; // 替换为实际存在的订单ID
Order order = orderMapper.selectById(orderId);
if (order != null) {
log.info("查询到订单:订单ID={}, 用户ID={}, 订单编号={}, 金额={}",
order.getOrderId(), order.getUserId(), order.getOrderNo(), order.getTotalAmount());
// 查询订单明细
List<OrderItem> orderItems = orderItemMapper.selectByOrderId(orderId);
log.info("订单明细数量:{}", orderItems.size());
for (OrderItem item : orderItems) {
log.info("明细:商品ID={}, 商品名称={}, 数量={}, 单价={}",
item.getProductId(), item.getProductName(), item.getQuantity(), item.getPrice());
}
} else {
log.warn("未查询到订单,订单ID:{}", orderId);
}
log.info("===== 订单查询测试完成 =====");
}
/**
* 测试用户订单查询
*/
@Test
public void testSelectUserOrders() {
log.info("===== 开始测试用户订单查询 =====");
Long userId = 1L; // 查询用户ID为1的订单
List<Order> orders = orderService.getOrdersByUserId(userId);
log.info("用户{}的订单数量:{}", userId, orders.size());
for (Order order : orders) {
log.info("订单:订单ID={}, 订单编号={}, 金额={}, 状态={}",
order.getOrderId(), order.getOrderNo(), order.getTotalAmount(), order.getStatus());
}
log.info("===== 用户订单查询测试完成 =====");
}
/**
* 测试分片数据分布
*/
@Test
public void testShardingDistribution() {
log.info("===== 开始测试分片数据分布 =====");
Map<String, Object> statistics = orderService.getShardingStatistics();
log.info("订单总数:{}", statistics.get("totalOrders"));
@SuppressWarnings("unchecked")
List<Map<String, Object>> shardCounts =
(List<Map<String, Object>>) statistics.get("shardCounts");
log.info("各分片订单分布:");
for (Map<String, Object> shardInfo : shardCounts) {
log.info("分片:{}, 数量:{}", shardInfo.get("shard"), shardInfo.get("count"));
}
log.info("===== 分片数据分布测试完成 =====");
}
/**
* 测试订单状态更新
*/
@Test
public void testUpdateOrderStatus() {
log.info("===== 开始测试订单状态更新 =====");
Long orderId = 1234567890123456789L; // 替换为实际存在的订单ID
Integer newStatus = 1; // 更新为已支付状态
boolean success = orderService.updateOrderStatus(orderId, newStatus);
log.info("订单状态更新结果:{}, 订单ID:{}, 新状态:{}", success, orderId, newStatus);
// 验证更新结果
Order updatedOrder = orderMapper.selectById(orderId);
if (updatedOrder != null) {
log.info("更新后订单状态:{}", updatedOrder.getStatus());
}
log.info("===== 订单状态更新测试完成 =====");
}
/**
* 测试跨分片查询
*/
@Test
public void testCrossShardQuery() {
log.info("===== 开始测试跨分片查询 =====");
// 查询所有订单(会跨多个分片)
List<Order> allOrders = orderMapper.selectList(null);
log.info("总订单数:{}", allOrders.size());
// 统计各用户的订单数
Map<Long, Long> userOrderCount = new HashMap<>();
for (Order order : allOrders) {
Long userId = order.getUserId();
userOrderCount.put(userId, userOrderCount.getOrDefault(userId, 0L) + 1);
}
log.info("各用户订单统计:");
for (Map.Entry<Long, Long> entry : userOrderCount.entrySet()) {
log.info("用户ID:{}, 订单数:{}", entry.getKey(), entry.getValue());
}
log.info("===== 跨分片查询测试完成 =====");
}
/**
* 测试批量插入
*/
@Test
public void testBatchInsert() {
log.info("===== 开始测试批量插入 =====");
List<Order> orders = new ArrayList<>();
// 准备批量数据
for (int i = 1; i <= 100; i++) {
Order order = new Order();
order.setUserId((long) (i % 5 + 1)); // 5个用户
order.setOrderNo("BATCH" + System.currentTimeMillis() + i);
order.setTotalAmount(new BigDecimal(String.valueOf(100 + i)));
order.setStatus(0);
order.setCreateTime(new Date());
order.setUpdateTime(new Date());
orders.add(order);
}
// 批量插入
long startTime = System.currentTimeMillis();
for (Order order : orders) {
orderMapper.insert(order);
}
long endTime = System.currentTimeMillis();
log.info("批量插入{}条订单完成,耗时:{}ms", orders.size(), (endTime – startTime));
log.info("===== 批量插入测试完成 =====");
}
/**
* 测试分片路由效果
*/
@Test
public void testShardingRouting() {
log.info("===== 开始测试分片路由效果 =====");
// 测试不同用户ID的路由结果
for (long userId = 1; userId <= 10; userId++) {
Order order = new Order();
order.setUserId(userId);
order.setOrderNo("ROUTE" + System.currentTimeMillis() + userId);
order.setTotalAmount(new BigDecimal("88.88"));
order.setStatus(0);
order.setCreateTime(new Date());
order.setUpdateTime(new Date());
orderMapper.insert(order);
// 计算路由结果
int dbIndex = (int) (userId % 2);
int tableIndex = (int) (order.getOrderId() % 2);
log.info("用户ID:{}, 订单ID:{}, 路由到:ds{}.t_order_{}",
userId, order.getOrderId(), dbIndex, tableIndex);
}
log.info("===== 分片路由效果测试完成 =====");
}
}
5.2 测试结果展示
订单插入测试结果:
===== 开始测试订单插入 =====
订单插入结果:true, 订单ID:1234567890123456789, 用户ID:1, 订单编号:ORD17047296000001
订单插入结果:true, 订单ID:1234567890123456790, 用户ID:2, 订单编号:ORD17047296000002
订单插入结果:true, 订单ID:1234567890123456791, 用户ID:1, 订单编号:ORD17047296000003
…
===== 订单插入测试完成 =====
分片数据分布统计结果:
===== 开始测试分片数据分布 =====
订单总数:1000
各分片订单分布:
分片:0, 数量:502
分片:1, 数量:498
===== 分片数据分布测试完成 =====
分片路由效果验证:
===== 开始测试分片路由效果 =====
用户ID:1, 订单ID:1234567890123456789, 路由到:ds1.t_order_1
用户ID:2, 订单ID:1234567890123456790, 路由到:ds0.t_order_0
用户ID:3, 订单ID:1234567890123456791, 路由到:ds1.t_order_1
用户ID:4, 订单ID:1234567890123456792, 路由到:ds0.t_order_0
…
===== 分片路由效果测试完成 =====
性能测试结果:
===== 开始测试批量插入 =====
批量插入100条订单完成,耗时:1523ms
平均每条订单插入耗时:15.23ms
===== 批量插入测试完成 =====
5.3 数据分片效果验证
分片均匀性验证:
— 验证各分片数据分布
SELECT
CONCAT('ds', MOD(user_id, 2)) AS database_name,
CONCAT('t_order_', MOD(order_id, 2)) AS table_name,
COUNT(*) AS count
FROM t_order
GROUP BY database_name, table_name
ORDER BY database_name, table_name;
预期结果:
| database_name | table_name | count |
|—————|————–|——-|
| ds0 | t_order_0 | 250 |
| ds0 | t_order_1 | 250 |
| ds1 | t_order_0 | 250 |
| ds1 | t_order_1 | 250 |
绑定表验证:
— 验证订单和明细在同一物理表中
SELECT
o.order_id,
o.order_no,
i.item_id,
i.product_name
FROM t_order o
JOIN t_order_item i ON o.order_id = i.order_id
WHERE o.order_id = 1234567890123456789;
6. 进阶优化
6.1 性能调优策略
6.1.1 分片键选择优化
分片键选择原则:
优化案例:
# 错误的分片键选择
database-strategy:
standard:
sharding-column: status # 状态字段只有几个值,会导致数据倾斜
# 正确的分片键选择
database-strategy:
standard:
sharding-column: user_id # 用户ID分布均匀,查询频率高
复合分片策略:
database-strategy:
complex:
sharding-columns: user_id,create_time
sharding-algorithm-name: complex–db–sharding
sharding-algorithms:
complex-db-sharding:
type: CLASS_BASED
props:
strategy-class-name: com.example.ComplexDatabaseShardingAlgorithm
6.1.2 索引优化策略
分片键索引:
— 分片键必须建立索引
CREATE INDEX idx_user_id ON t_order(user_id);
CREATE INDEX idx_order_id ON t_order(order_id);
全局索引设计: 对于非分片键的查询,可以考虑以下方案:
方案1:冗余索引表
— 创建冗余索引表
CREATE TABLE t_order_index (
order_id BIGINT PRIMARY KEY,
user_id BIGINT,
order_no VARCHAR(64),
create_time DATETIME,
INDEX idx_order_no(order_no)
);
方案2:Elasticsearch搜索引擎
// 将订单数据同步到Elasticsearch
@Document(indexName = "order_index")
public class OrderDocument {
@Id
private Long orderId;
@Field(type = FieldType.Keyword)
private String orderNo;
@Field(type = FieldType.Long)
private Long userId;
// 其他字段…
}
6.1.3 查询性能优化
避免全分片扫描:
// 错误:会导致全分片扫描
List<Order> orders = orderMapper.selectList(
new QueryWrapper<Order>().like("order_no", "ORD")
);
// 正确:添加分片键条件,限定查询范围
List<Order> orders = orderMapper.selectList(
new QueryWrapper<Order>()
.eq("user_id", userId)
.like("order_no", "ORD")
);
分页查询优化:
// 错误:深分页性能差
Page<Order> page = new Page<>(10000, 20);
// 正确:使用游标分页
Page<Order> page = new Page<>(1, 20);
QueryWrapper<Order> wrapper = new QueryWrapper<Order>()
.gt("order_id", lastOrderId)
.orderByAsc("order_id");
批量操作优化:
// 优化批量插入
@Transactional
public void batchInsertOrders(List<Order> orders) {
// 按分片分组
Map<Integer, List<Order>> groupedOrders = orders.stream()
.collect(Collectors.groupingBy(
order -> (int) (order.getUserId() % 2)
));
// 分批插入
for (Map.Entry<Integer, List<Order>> entry : groupedOrders.entrySet()) {
int batchSize = 100;
List<List<Order>> batches = Lists.partition(entry.getValue(), batchSize);
for (List<Order> batch : batches) {
for (Order order : batch) {
orderMapper.insert(order);
}
}
}
}
6.1.4 连接池调优
HikariCP优化配置:
spring:
shardingsphere:
datasource:
ds0:
hikari:
minimum-idle: 10 # 最小空闲连接数
maximum-pool-size: 50 # 最大连接池大小
connection-timeout: 30000 # 连接超时时间
idle-timeout: 600000 # 空闲连接超时时间
max-lifetime: 1800000 # 连接最大生命周期
connection-test-query: SELECT 1
leak-detection-threshold: 60000 # 连接泄漏检测阈值
连接池监控:
@Component
public class HikariPoolMonitor {
@Autowired
private DataSource dataSource;
@Scheduled(fixedRate = 60000) // 每分钟监控一次
public void monitorHikariPool() {
if (dataSource instanceof HikariDataSource) {
HikariPoolMXBean poolProxy = ((HikariDataSource) dataSource).getHikariPoolMXBean();
log.info("HikariCP连接池状态:");
log.info("活跃连接数:{}", poolProxy.getActiveConnections());
log.info("空闲连接数:{}", poolProxy.getIdleConnections());
log.info("总连接数:{}", poolProxy.getTotalConnections());
log.info("线程等待连接数:{}", poolProxy.getThreadsAwaitingConnection());
}
}
}
6.2 常见问题解决方案
6.2.1 数据倾斜问题
问题现象:
- 某些分片数据量远大于其他分片
- 某些数据库负载过高,响应时间过长
解决方案:
方案1:优化分片算法
public class ConsistentHashingShardingAlgorithm implements PreciseShardingAlgorithm<Long> {
private final TreeMap<Long, String> virtualNodes = new TreeMap<>();
private static final int VIRTUAL_NODE_COUNT = 150; // 虚拟节点数量
public ConsistentHashingShardingAlgorithm(List<String> targetNames) {
for (String targetName : targetNames) {
for (int i = 0; i < VIRTUAL_NODE_COUNT; i++) {
long hash = hash(targetName + ":" + i);
virtualNodes.put(hash, targetName);
}
}
}
@Override
public String doSharding(Collection<String> availableTargetNames,
PreciseShardingValue<Long> shardingValue) {
Long value = shardingValue.getValue();
long hash = hash(value.toString());
Map.Entry<Long, String> entry = virtualNodes.ceilingEntry(hash);
if (entry == null) {
entry = virtualNodes.firstEntry();
}
return entry.getValue();
}
private long hash(String key) {
// 使用一致性哈希算法
return Hashing.consistentHash(key.hashCode(), Integer.MAX_VALUE);
}
}
方案2:动态扩容
# 使用ShardingSphere-Scaling进行数据迁移
rules:
scaling:
tables:
t_order:
input-data-sources: ds0,ds1
output-data-sources: ds0,ds1,ds2,ds3
6.2.2 跨分片JOIN问题
问题现象:
- 订单和订单明细跨分片存储
- JOIN查询性能差或无法执行
- 跨库查询导致全表扫描,响应时间过长
解决方案:
方案1:绑定表(推荐,适用于分片规则一致的场景)
核心原理: 绑定表是指分片规则完全一致的主表和子表。ShardingSphere在处理绑定表时,会根据分片键将相关数据路由到同一个物理分片中,从而避免跨分片JOIN。
YAML配置:
spring:
shardingsphere:
rules:
sharding:
# 绑定表配置
binding-tables:
– t_order,t_order_item # 订单表和订单明细表为绑定表
实现要求:
SQL示例:
// 在Service层实现绑定表JOIN查询
@Service
public class OrderServiceImpl implements OrderService {
/**
* 查询订单及其明细(绑定表优化)
* 说明:由于t_order和t_order_item配置为绑定表,
* ShardingSphere会自动将它们路由到同一个物理分片
*/
public Map<String, Object> getOrderWithItems(Long orderId) {
// 第一步:查询订单信息(会路由到具体分片)
Order order = orderMapper.selectById(orderId);
// 第二步:查询订单明细(由于是绑定表,会路由到同一分片)
List<OrderItem> orderItems = orderItemMapper.selectByOrderId(orderId);
// 第三步:组装返回结果
Map<String, Object> result = new HashMap<>();
result.put("order", order);
result.put("orderItems", orderItems);
result.put("totalItems", orderItems.size());
result.put("totalAmount", order.getTotalAmount());
return result;
}
}
Mapper实现:
@Mapper
public interface OrderMapper extends BaseMapper<Order> {
// 使用MyBatis-Plus提供的通用方法
// selectById会自动根据order_id分片键进行路由
}
@Mapper
public interface OrderItemMapper extends BaseMapper<OrderItem> {
/**
* 根据订单ID查询订单明细
* 说明:由于配置了绑定表,此查询会自动路由到与订单表相同的分片
*/
@Select("SELECT * FROM t_order_item WHERE order_id = #{orderId}")
List<OrderItem> selectByOrderId(Long orderId);
}
性能优势:
- 避免了跨分片数据传输
- 在单个分片内完成JOIN,性能接近单表查询
- 充分利用数据库本地JOIN优化
方案2:应用层关联(适用于查询条件明确、结果集小的场景)
核心原理: 将复杂的SQL JOIN拆分为多个简单的单表查询,在应用层进行数据组装。适用于查询结果集较小、可以通过分片键限定查询范围的场景。
完整实现代码:
package com.example.shardingsphere.service.impl;
import com.example.shardingsphere.entity.Order;
import com.example.shardingsphere.entity.OrderItem;
import com.example.shardingsphere.mapper.OrderItemMapper;
import com.example.shardingsphere.mapper.OrderMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.util.*;
import java.util.stream.Collectors;
/**
* 应用层关联查询服务
*/
@Slf4j
@Service
public class ApplicationLayerJoinService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private OrderItemMapper orderItemMapper;
/**
* 场景1:查询用户订单及其明细(通过用户ID限定分片范围)
*
* @param userId 用户ID(分库键)
* @return 用户订单列表及明细
*/
public List<Map<String, Object>> getUserOrdersWithItems(Long userId) {
log.info("查询用户订单及明细,用户ID:{}", userId);
// 第一步:查询用户的所有订单
List<Order> orders = orderMapper.selectList(
new com.baomidou.mybatisplus.core.conditions.query.QueryWrapper<Order>()
.eq("user_id", userId)
.orderByDesc("create_time")
);
if (orders.isEmpty()) {
log.info("用户{}暂无订单", userId);
return Collections.emptyList();
}
log.info("用户{}共有{}条订单", userId, orders.size());
// 第二步:提取订单ID列表
List<Long> orderIds = orders.stream()
.map(Order::getOrderId)
.collect(Collectors.toList());
// 第三步:批量查询订单明细
// 说明:由于订单明细表的分片键是order_id,这里需要跨分片查询
// 但通过订单ID列表,可以减少查询次数
List<OrderItem> allOrderItems = orderItemMapper.selectList(
new com.baomidou.mybatisplus.core.conditions.query.QueryWrapper<OrderItem>()
.in("order_id", orderIds)
);
// 第四步:按订单ID分组明细
Map<Long, List<OrderItem>> orderItemMap = allOrderItems.stream()
.collect(Collectors.groupingBy(OrderItem::getOrderId));
// 第五步:组装结果
List<Map<String, Object>> result = new ArrayList<>();
for (Order order : orders) {
Map<String, Object> orderData = new HashMap<>();
orderData.put("order", order);
// 获取该订单的明细
List<OrderItem> items = orderItemMap.getOrDefault(order.getOrderId(), Collections.emptyList());
orderData.put("orderItems", items);
orderData.put("itemCount", items.size());
result.add(orderData);
}
log.info("查询完成,返回{}条订单数据", result.size());
return result;
}
/**
* 场景2:查询订单及明细(单个订单)
*
* @param orderId 订单ID
* @return 订单及明细信息
*/
public Map<String, Object> getOrderWithItems(Long orderId) {
log.info("查询单个订单及明细,订单ID:{}", orderId);
// 第一步:查询订单
Order order = orderMapper.selectById(orderId);
if (order == null) {
log.warn("订单不存在,订单ID:{}", orderId);
return null;
}
// 第二步:查询订单明细
List<OrderItem> orderItems = orderItemMapper.selectByOrderId(orderId);
// 第三步:组装结果
Map<String, Object> result = new HashMap<>();
result.put("order", order);
result.put("orderItems", orderItems);
result.put("itemCount", orderItems.size());
log.info("查询完成,订单{}包含{}条明细", orderId, orderItems.size());
return result;
}
/**
* 场景3:分页查询用户订单(避免内存分页)
*
* @param userId 用户ID
* @param pageNum 页码
* @param pageSize 每页大小
* @return 分页结果
*/
public Map<String, Object> getUserOrdersPage(Long userId, int pageNum, int pageSize) {
log.info("分页查询用户订单,用户ID:{},页码:{},每页:{}", userId, pageNum, pageSize);
// 计算分页参数
int offset = (pageNum – 1) * pageSize;
// 第一步:分页查询订单
List<Order> orders = orderMapper.selectList(
new com.baomidou.mybatisplus.core.conditions.query.QueryWrapper<Order>()
.eq("user_id", userId)
.orderByDesc("create_time")
.last("LIMIT " + offset + ", " + pageSize)
);
// 第二步:查询总记录数
int total = orderMapper.selectCount(
new com.baomidou.mybatisplus.core.conditions.query.QueryWrapper<Order>()
.eq("user_id", userId)
);
// 第三步:组装分页结果
Map<String, Object> pageResult = new HashMap<>();
pageResult.put("list", orders);
pageResult.put("total", total);
pageResult.put("pageNum", pageNum);
pageResult.put("pageSize", pageSize);
pageResult.put("pages", (total + pageSize – 1) / pageSize);
log.info("分页查询完成,总记录数:{},当前页记录数:{}", total, orders.size());
return pageResult;
}
/**
* 场景4:多表关联查询(订单+用户+商品)
* 说明:演示如何通过应用层关联多个表
*
* @param orderId 订单ID
* @return 完整的订单信息
*/
public Map<String, Object> getCompleteOrderInfo(Long orderId) {
log.info("查询完整订单信息,订单ID:{}", orderId);
// 第一步:查询订单信息
Order order = orderMapper.selectById(orderId);
if (order == null) {
return null;
}
// 第二步:查询订单明细
List<OrderItem> orderItems = orderItemMapper.selectByOrderId(orderId);
// 第三步:这里可以调用其他Service查询用户信息、商品信息等
// UserInfo userInfo = userService.getUserInfo(order.getUserId());
// List<ProductInfo> products = productService.getProductInfo(orderItems);
// 第四步:组装完整结果
Map<String, Object> result = new HashMap<>();
result.put("order", order);
result.put("orderItems", orderItems);
// result.put("userInfo", userInfo);
// result.put("products", products);
return result;
}
}
性能优化技巧:
/**
* 批量查询优化
* 说明:避免N+1查询问题,使用批量查询代替循环查询
*/
public List<Order> getOrdersOptimization(List<Long> orderIds) {
// 错误示范:循环查询(N+1问题)
// List<Order> result = new ArrayList<>();
// for (Long orderId : orderIds) {
// Order order = orderMapper.selectById(orderId); // 每次循环都执行一次查询
// result.add(order);
// }
// 正确示范:批量查询
List<Order> result = orderMapper.selectBatchIds(orderIds);
return result;
}
/**
* 并行查询优化
* 说明:对于没有依赖关系的查询,可以使用并行执行
*/
public Map<String, Object> parallelQuery(Long orderId) {
// 使用CompletableFuture并行查询
CompletableFuture<Order> orderFuture = CompletableFuture.supplyAsync(() ->
orderMapper.selectById(orderId)
);
CompletableFuture<List<OrderItem>> itemsFuture = CompletableFuture.supplyAsync(() ->
orderItemMapper.selectByOrderId(orderId)
);
// 等待所有查询完成
CompletableFuture.allOf(orderFuture, itemsFuture).join();
// 组装结果
Map<String, Object> result = new HashMap<>();
result.put("order", orderFuture.join());
result.put("orderItems", itemsFuture.join());
return result;
}
方案3:数据冗余(适用于读多写少、实时性要求不高的场景)
核心原理: 在子表中冗余主表的分片键,或者在主表中冗余子表的聚合信息,将跨表查询转换为单表查询。
实现方案1:订单明细表冗余用户ID
— 步骤1:在订单明细表中添加用户ID字段
ALTER TABLE t_order_item ADD COLUMN user_id BIGINT COMMENT '用户ID(冗余)';
— 步骤2:为冗余字段建立索引
CREATE INDEX idx_user_id ON t_order_item(user_id);
— 步骤3:修改数据同步逻辑(在插入/更新订单明细时同步用户ID)
— 这一步需要在应用层完成,或者使用数据库触发器
Java代码实现:
@Service
public class OrderItemService {
@Autowired
private OrderItemMapper orderItemMapper;
@Autowired
private OrderMapper orderMapper;
/**
* 插入订单明细(带用户ID冗余)
*/
@Transactional(rollbackFor = Exception.class)
public void insertOrderItemWithRedundancy(OrderItem orderItem) {
// 查询订单信息获取用户ID
Order order = orderMapper.selectById(orderItem.getOrderId());
if (order != null) {
// 冗余用户ID到订单明细表
orderItem.setUserId(order.getUserId());
}
// 插入订单明细
orderItemMapper.insert(orderItem);
}
/**
* 查询用户订单明细(利用冗余字段直接分片)
* 说明:由于订单明细表有了user_id字段,可以直接按用户ID分片查询
*/
public List<OrderItem> getUserOrderItems(Long userId) {
return orderItemMapper.selectList(
new com.baomidou.mybatisplus.core.conditions.query.QueryWrapper<OrderItem>()
.eq("user_id", userId)
.orderByDesc("create_time")
);
}
}
实现方案2:订单表冗余订单摘要信息
— 在订单表中添加冗余字段
ALTER TABLE t_order ADD COLUMN item_count INT DEFAULT 0 COMMENT '商品数量(冗余)';
ALTER TABLE t_order ADD COLUMN item_summary TEXT COMMENT '订单明细摘要(冗余)';
Java代码实现:
@Service
public class OrderServiceWithRedundancy {
@Autowired
private OrderMapper orderMapper;
@Autowired
private OrderItemMapper orderItemMapper;
/**
* 创建订单(更新冗余字段)
*/
@Transactional(rollbackFor = Exception.class)
public void createOrderWithRedundancy(Order order, List<OrderItem> orderItems) {
// 插入订单
orderMapper.insert(order);
// 插入订单明细
for (OrderItem item : orderItems) {
item.setOrderId(order.getOrderId());
orderItemMapper.insert(item);
}
// 更新订单表的冗余字段
Order updateOrder = new Order();
updateOrder.setOrderId(order.getOrderId());
updateOrder.setItemCount(orderItems.size());
// 生成订单明细摘要
String summary = generateOrderSummary(orderItems);
updateOrder.setItemSummary(summary);
orderMapper.updateById(updateOrder);
}
/**
* 生成订单摘要
*/
private String generateOrderSummary(List<OrderItem> orderItems) {
StringBuilder summary = new StringBuilder();
for (OrderItem item : orderItems) {
if (summary.length() > 0) {
summary.append(", ");
}
summary.append(item.getProductName())
.append(" x ")
.append(item.getQuantity());
}
return summary.toString();
}
/**
* 查询用户订单摘要(利用冗余字段,避免JOIN订单明细表)
*/
public List<Order> getUserOrderSummaries(Long userId) {
return orderMapper.selectList(
new com.baomidou.mybatisplus.core.conditions.query.QueryWrapper<Order>()
.eq("user_id", userId)
.select("order_id", "order_no", "total_amount", "status",
"create_time", "item_count", "item_summary")
.orderByDesc("create_time")
);
}
}
方案4:Elasticsearch搜索引擎(适用于全文检索、复杂查询场景)
核心原理: 将分库分表的数据同步到Elasticsearch,利用Elasticsearch强大的全文检索和聚合查询能力,解决跨分片JOIN问题。
实现步骤:
步骤1:引入Elasticsearch依赖
<!– Elasticsearch依赖 –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-elasticsearch</artifactId>
</dependency>
<!– Elasticsearch高级客户端 –>
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
</dependency>
步骤2:定义Elasticsearch文档
package com.example.shardingsphere.elasticsearch;
import lombok.Data;
import org.springframework.data.annotation.Id;
import org.springframework.data.elasticsearch.annotations.Document;
import org.springframework.data.elasticsearch.annotations.Field;
import org.springframework.data.elasticsearch.annotations.FieldType;
import java.math.BigDecimal;
import java.util.Date;
import java.util.List;
/**
* 订单Elasticsearch文档
*/
@Data
@Document(indexName = "order_index", createIndex = true)
public class OrderDocument {
@Id
private Long orderId;
@Field(type = FieldType.Long)
private Long userId;
@Field(type = FieldType.Keyword)
private String orderNo;
@Field(type = FieldType.Double)
private BigDecimal totalAmount;
@Field(type = FieldType.Integer)
private Integer status;
@Field(type = FieldType.Date)
private Date createTime;
@Field(type = FieldType.Nested)
private List<OrderItemDocument> orderItems;
/**
* 订单明细文档
*/
@Data
public static class OrderItemDocument {
@Field(type = FieldType.Long)
private Long itemId;
@Field(type = FieldType.Long)
private Long productId;
@Field(type = FieldType.Text, analyzer = "ik_max_word")
private String productName;
@Field(type = FieldType.Integer)
private Integer quantity;
@Field(type = FieldType.Double)
private BigDecimal price;
}
}
步骤3:实现数据同步
package com.example.shardingsphere.elasticsearch;
import com.example.shardingsphere.entity.Order;
import com.example.shardingsphere.entity.OrderItem;
import com.example.shardingsphere.mapper.OrderItemMapper;
import com.example.shardingsphere.mapper.OrderMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.elasticsearch.core.ElasticsearchRestTemplate;
import org.springframework.data.elasticsearch.core.query.IndexQuery;
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.stream.Collectors;
/**
* Elasticsearch数据同步服务
*/
@Slf4j
@Service
public class OrderElasticsearchService {
@Autowired
private ElasticsearchRestTemplate elasticsearchTemplate;
@Autowired
private OrderMapper orderMapper;
@Autowired
private OrderItemMapper orderItemMapper;
/**
* 同步单个订单到Elasticsearch
*/
public void syncOrderToEs(Long orderId) {
log.info("开始同步订单到Elasticsearch,订单ID:{}", orderId);
try {
// 查询订单信息
Order order = orderMapper.selectById(orderId);
if (order == null) {
log.warn("订单不存在,订单ID:{}", orderId);
return;
}
// 查询订单明细
List<OrderItem> orderItems = orderItemMapper.selectByOrderId(orderId);
// 转换为Elasticsearch文档
OrderDocument document = convertToDocument(order, orderItems);
// 索引到Elasticsearch
IndexQuery indexQuery = new IndexQuery();
indexQuery.setId(orderId.toString());
indexQuery.setObject(document);
elasticsearchTemplate.index(indexQuery);
log.info("订单同步到Elasticsearch成功,订单ID:{}", orderId);
} catch (Exception e) {
log.error("同步订单到Elasticsearch失败,订单ID:{}", orderId, e);
}
}
/**
* 批量同步订单到Elasticsearch
*/
public void batchSyncOrdersToEs(List<Long> orderIds) {
log.info("开始批量同步订单到Elasticsearch,订单数量:{}", orderIds.size());
List<IndexQuery> indexQueries = orderIds.stream()
.map(orderId -> {
Order order = orderMapper.selectById(orderId);
if (order == null) {
return null;
}
List<OrderItem> orderItems = orderItemMapper.selectByOrderId(orderId);
OrderDocument document = convertToDocument(order, orderItems);
IndexQuery indexQuery = new IndexQuery();
indexQuery.setId(orderId.toString());
indexQuery.setObject(document);
return indexQuery;
})
.filter(Objects::nonNull)
.collect(Collectors.toList());
elasticsearchTemplate.bulkIndex(indexQueries);
log.info("批量同步订单到Elasticsearch完成,成功数量:{}", indexQueries.size());
}
/**
* 转换为Elasticsearch文档
*/
private OrderDocument convertToDocument(Order order, List<OrderItem> orderItems) {
OrderDocument document = new OrderDocument();
document.setOrderId(order.getOrderId());
document.setUserId(order.getUserId());
document.setOrderNo(order.getOrderNo());
document.setTotalAmount(order.getTotalAmount());
document.setStatus(order.getStatus());
document.setCreateTime(order.getCreateTime());
// 转换订单明细
List<OrderDocument.OrderItemDocument> itemDocuments = orderItems.stream()
.map(item -> {
OrderDocument.OrderItemDocument itemDoc = new OrderDocument.OrderItemDocument();
itemDoc.setItemId(item.getItemId());
itemDoc.setProductId(item.getProductId());
itemDoc.setProductName(item.getProductName());
itemDoc.setQuantity(item.getQuantity());
itemDoc.setPrice(item.getPrice());
return itemDoc;
})
.collect(Collectors.toList());
document.setOrderItems(itemDocuments);
return document;
}
}
步骤4:实现复杂查询
package com.example.shardingsphere.elasticsearch;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.index.query.BoolQueryBuilder;
import org.elasticsearch.index.query.QueryBuilders;
import org.elasticsearch.search.aggregations.AggregationBuilders;
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
import org.elasticsearch.search.sort.SortBuilders;
import org.elasticsearch.search.sort.SortOrder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.domain.PageRequest;
import org.springframework.data.elasticsearch.core.ElasticsearchRestTemplate;
import org.springframework.data.elasticsearch.core.NativeSearchQueryBuilder;
import org.springframework.data.elasticsearch.core.SearchHit;
import org.springframework.data.elasticsearch.core.SearchHits;
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
/**
* Elasticsearch查询服务
*/
@Slf4j
@Service
public class OrderEsQueryService {
@Autowired
private ElasticsearchRestTemplate elasticsearchTemplate;
/**
* 全文检索订单(支持商品名称搜索)
* 说明:可以搜索商品名称,解决跨表查询问题
*/
public List<OrderDocument> searchOrdersByProduct(String productName) {
log.info("全文检索订单,商品名称:{}", productName);
NativeSearchQueryBuilder queryBuilder = new NativeSearchQueryBuilder();
// 使用嵌套查询搜索商品名称
BoolQueryBuilder boolQuery = QueryBuilders.boolQuery();
boolQuery.must(QueryBuilders.nestedQuery(
"orderItems",
QueryBuilders.matchQuery("orderItems.productName", productName),
org.elasticsearch.index.query.ScoreMode.None
));
queryBuilder.withQuery(boolQuery);
queryBuilder.withSort(SortBuilders.fieldSort("createTime").order(SortOrder.DESC));
queryBuilder.withPageable(PageRequest.of(0, 20));
SearchHits<OrderDocument> searchHits = elasticsearchTemplate.search(queryBuilder.build(), OrderDocument.class);
log.info("全文检索完成,命中数量:{}", searchHits.getTotalHits());
return searchHits.getSearchHits().stream()
.map(SearchHit::getContent)
.collect(java.util.stream.Collectors.toList());
}
/**
* 复杂条件查询
*/
public List<OrderDocument> complexSearch(Long userId, String productName, Integer status) {
log.info("复杂条件查询,用户ID:{},商品名称:{},状态:{}", userId, productName, status);
BoolQueryBuilder boolQuery = QueryBuilders.boolQuery();
// 用户ID查询
if (userId != null) {
boolQuery.must(QueryBuilders.termQuery("userId", userId));
}
// 商品名称查询(嵌套)
if (productName != null && !productName.isEmpty()) {
boolQuery.must(QueryBuilders.nestedQuery(
"orderItems",
QueryBuilders.matchQuery("orderItems.productName", productName),
org.elasticsearch.index.query.ScoreMode.None
));
}
// 订单状态查询
if (status != null) {
boolQuery.must(QueryBuilders.termQuery("status", status));
}
NativeSearchQueryBuilder queryBuilder = new NativeSearchQueryBuilder()
.withQuery(boolQuery)
.withSort(SortBuilders.fieldSort("createTime").order(SortOrder.DESC))
.withPageable(PageRequest.of(0, 20));
SearchHits<OrderDocument> searchHits = elasticsearchTemplate.search(queryBuilder.build(), OrderDocument.class);
return searchHits.getSearchHits().stream()
.map(SearchHit::getContent)
.collect(java.util.stream.Collectors.toList());
}
/**
* 聚合查询统计
*/
public Map<String, Object> aggregationStatistics(Long userId) {
log.info("执行聚合查询统计,用户ID:{}", userId);
NativeSearchQueryBuilder queryBuilder = new NativeSearchQueryBuilder();
// 查询条件
BoolQueryBuilder boolQuery = QueryBuilders.boolQuery();
boolQuery.must(QueryBuilders.termQuery("userId", userId));
queryBuilder.withQuery(boolQuery);
// 订单状态聚合
queryBuilder.addAggregation(
AggregationBuilders.terms("status_count").field("status")
);
// 订单商品聚合
queryBuilder.addAggregation(
AggregationBuilders.nested("product_count", "orderItems")
.subAggregation(AggregationBuilders.terms("product_name").field("orderItems.productName"))
);
SearchHits<OrderDocument> searchHits = elasticsearchTemplate.search(queryBuilder.build(), OrderDocument.class);
// 提取聚合结果
Terms statusTerms = searchHits.getAggregations().get("status_count");
Terms productTerms = searchHits.getAggregations().get("product_count");
Map<String, Object> result = new java.util.HashMap<>();
result.put("totalOrders", searchHits.getTotalHits());
result.put("statusDistribution", statusTerms.getBuckets());
result.put("topProducts", productTerms);
return result;
}
}
方案5:全局索引表(适用于需要频繁按非分片键查询的场景)
核心原理: 创建专门的全局索引表,存储逻辑表的全局索引信息,通过索引表快速定位数据所在的物理分片。
实现方案:
— 创建全局索引表
CREATE TABLE t_order_global_index (
order_id BIGINT PRIMARY KEY COMMENT '订单ID',
order_no VARCHAR(64) UNIQUE KEY COMMENT '订单编号',
user_id BIGINT NOT NULL COMMENT '用户ID',
db_shard INT NOT NULL COMMENT '数据库分片索引',
table_shard INT NOT NULL COMMENT '表分片索引',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
INDEX idx_order_no(order_no),
INDEX idx_user_id(user_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单全局索引表';
Java代码实现:
@Service
public class OrderGlobalIndexService {
@Autowired
private JdbcTemplate jdbcTemplate;
/**
* 创建全局索引
*/
@Transactional(rollbackFor = Exception.class)
public void createGlobalIndex(Order order) {
// 计算分片索引
int dbShard = (int) (order.getUserId() % 2);
int tableShard = (int) (order.getOrderId() % 2);
// 插入全局索引
String sql = "INSERT INTO t_order_global_index (order_id, order_no, user_id, db_shard, table_shard) " +
"VALUES (?, ?, ?, ?, ?)";
jdbcTemplate.update(sql,
order.getOrderId(),
order.getOrderNo(),
order.getUserId(),
dbShard,
tableShard
);
}
/**
* 通过订单编号查询分片信息
*/
public Map<String, Object> getShardInfoByOrderNo(String orderNo) {
String sql = "SELECT * FROM t_order_global_index WHERE order_no = ?";
return jdbcTemplate.queryForMap(sql, orderNo);
}
/**
* 通过全局索引查询订单
*/
public Order getOrderByOrderNo(String orderNo) {
// 第一步:查询全局索引
Map<String, Object> indexInfo = getShardInfoByOrderNo(orderNo);
if (indexInfo == null) {
return null;
}
// 第二步:根据分片信息路由查询
Long orderId = (Long) indexInfo.get("order_id");
return orderMapper.selectById(orderId);
}
}
6.2.3 分布式事务问题
问题现象:
- 跨库操作数据不一致
- 事务回滚失败
- 性能问题严重
解决方案:
方案1:使用XA事务
@Transactional
@ShardingTransactionType(TransactionType.XA)
public void createOrderWithXA(Order order, List<OrderItem> orderItems) {
// 强一致性事务
orderMapper.insert(order);
for (OrderItem item : orderItems) {
orderItemMapper.insert(item);
}
}
方案2:使用SEATA事务
@GlobalTransactional // Seata全局事务
public void createOrderWithSeata(Order order, List<OrderItem> orderItems) {
// 最终一致性事务
orderMapper.insert(order);
for (OrderItem item : orderItems) {
orderItemMapper.insert(item);
}
}
方案3:业务补偿
public void createOrderWithCompensation(Order order, List<OrderItem> orderItems) {
try {
// 执行主业务逻辑
orderMapper.insert(order);
for (OrderItem item : orderItems) {
orderItemMapper.insert(item);
}
} catch (Exception e) {
// 执行补偿操作
compensateOrder(order.getOrderId());
throw new RuntimeException("创建订单失败,已执行补偿操作");
}
}
private void compensateOrder(Long orderId) {
// 删除已插入的订单
orderMapper.deleteById(orderId);
// 删除已插入的订单明细
orderItemMapper.delete(
new QueryWrapper<OrderItem>().eq("order_id", orderId)
);
}
6.2.4 分布式ID问题
问题现象:
- 主键冲突
- ID不连续
- 性能问题
解决方案:
方案1:Snowflake算法
spring:
shardingsphere:
rules:
sharding:
key-generators:
snowflake:
type: SNOWFLAKE
props:
worker-id: 123 # 工作机器ID
max-tolerate-time-difference-milliseconds: 1000 # 最大容忍时钟回差
方案2:数据库自增
— 为每个分片设置不同的自增起始值
— ds0数据库
ALTER TABLE t_order_0 AUTO_INCREMENT = 1;
ALTER TABLE t_order_1 AUTO_INCREMENT = 1000000;
— ds1数据库
ALTER TABLE t_order_0 AUTO_INCREMENT = 2000000;
ALTER TABLE t_order_1 AUTO_INCREMENT = 3000000;
方案3:Redis生成ID
@Service
public class RedisIdGenerator {
@Autowired
private StringRedisTemplate redisTemplate;
public Long generateId(String key) {
return redisTemplate.opsForValue().increment(key);
}
}
6.3 最佳实践总结
6.3.1 架构设计最佳实践
分片设计原则:
- 分片键选择:优先考虑查询频率和数据均匀性
- 分片数量:2的幂次方,便于扩容
- 预留容量:考虑未来3-5年的业务增长
事务处理原则:
- 尽量将事务控制在单一分片内
- 对于跨分片操作,根据业务一致性要求选择合适的事务类型
- 避免大事务,长事务
查询优化原则:
- 查询条件中包含分片键
- 避免深分页查询
- 合理使用缓存
6.3.2 运维管理最佳实践
监控告警:
- 监控各分片的数据量和访问量
- 监控数据库连接池状态
- 监控查询响应时间
- 设置合理的告警阈值
备份恢复:
- 制定完善的备份策略
- 定期测试备份数据的可用性
- 制定详细的恢复流程
扩容方案:
- 使用ShardingSphere-Scaling工具进行数据迁移
- 选择业务低峰期进行扩容
- 准备回滚方案
性能优化:
- 定期分析慢查询日志
- 优化索引和分片策略
- 监控和调整连接池配置
6.3.3 开发规范
代码规范:
- 避免在代码中硬编码物理表名
- 使用逻辑表名进行操作
- 合理使用事务注解
测试规范:
- 编写完整的单元测试
- 进行压力测试
- 验证分片路由的正确性
文档规范:
- 详细记录分片策略和规则
- 维护数据库设计文档
- 编写运维手册
7. 总结展望
7.1 项目经验总结
通过本项目的实战开发,我们总结了以下关键经验:
技术选型经验:
架构设计经验:
性能优化经验:
问题解决经验:
7.2 ShardingSphere适用场景分析
适用场景:
海量数据存储:
- 单表数据量超过千万级
- 需要进行水平扩展
- 示例:电商订单、用户行为日志
高并发访问:
- QPS超过单库承载能力
- 需要读写分离
- 示例:秒杀系统、即时通讯
复杂查询需求:
- 需要跨库JOIN
- 需要全局索引
- 示例:数据分析、报表系统
异构数据库整合:
- 需要整合多种数据库类型
- 需要透明的数据访问
- 示例:企业数据中台
不适用场景:
小规模数据:
- 数据量较小,单库单表足够
- 分库分表会带来额外的复杂度
简单查询:
- 查询模式简单,不需要复杂的路由
- 分库分表收益不明显
强实时性要求:
- 对数据实时性要求极高
- 分布式事务可能成为瓶颈
7.3 未来发展趋势
技术发展趋势:
云原生架构:
- 更好的云原生支持
- 容器化部署优化
- Kubernetes集成增强
智能化运维:
- 自动化的分片策略推荐
- 智能化的扩容缩容
- 基于AI的性能优化
多模态数据库支持:
- 更丰富的数据库类型支持
- 跨数据库查询优化
- 数据联邦能力增强
生态集成:
- 与大数据平台集成
- 支持流式计算
- 数据湖技术整合
应用场景拓展:
分布式数据库即服务:
- 提供更加便捷的DBaaS服务
- 降低使用门槛
- 完善的监控和运维体系
金融级分布式数据库:
- 增强分布式事务能力
- 提供更高的一致性保证
- 满足金融行业的合规要求
物联网数据处理:
- 支持海量时序数据
- 优化边缘计算场景
- 实时数据处理能力
7.4 学习建议
初级阶段:
中级阶段:
高级阶段:
学习资源:
- 官方文档:https://shardingsphere.apache.org/
- GitHub仓库:https://github.com/apache/shardingsphere
- 社区论坛:Apache ShardingSphere邮件列表
- 技术博客:关注社区技术专家的实践经验
结语
ShardingSphere作为一款优秀的分布式数据库中间件,为解决海量数据存储和高并发访问提供了完善的解决方案。通过本文的实战讲解,相信读者已经掌握了ShardingSphere的核心技术和应用方法。
在实际项目中,需要根据业务特点和需求,合理设计分片策略,选择合适的接入方式,并注重性能优化和问题解决。同时,要持续关注ShardingSphere的技术发展,不断学习和实践,才能在分布式数据库领域不断提升自己的技术能力。
希望本文能够帮助读者在实际项目中成功应用ShardingSphere,为企业的数字化转型提供有力的技术支撑。
相关阅读:
- ShardingSphere官方文档
- 分布式数据库技术选型指南
- 高并发系统架构设计实战




