ShardingSphere 实战
概述
ShardingSphere-JDBC 作为轻量级 Java 框架,在 JDBC 层提供分库分表能力。本文聚焦分库分表设计中的核心挑战,提供可落地的实战方案。
核心流程
SQL 进入
↓
ShardingPreparedStatement.execute()
↓
SQLRouteEngine.route() → 确定目标数据源 + 表
↓
SQLRewriteEngine.rewrite() → 改写 SQL(逻辑表→物理表)
↓
SQLExecutionEngine.execute() → 并行执行到多个数据源
↓
MergeEngine.merge() → 结果归并(分页/排序/聚合)
↓
返回结果一、分库分表设计与实现
1.1 核心概念
| 概念 | 说明 | 示例 |
|---|---|---|
| 逻辑表 | 应用视角的表名 | t_order |
| 物理表 | 实际存储的表名 | t_order_0, t_order_1 |
| 数据节点 | 映射关系 | ds0.t_order_0, ds1.t_order_1 |
| 分片键 | 分片路由字段 | order_id, user_id |
| 分片算法 | 路由到哪个分片 | order_id % 2 |
1.2 电商订单分片设计
yaml
spring:
shardingsphere:
datasource:
names: ds0, ds1, ds2, ds3 # 4 个数据库实例
ds0:
url: jdbc:mysql://192.168.1.10:3306/order_0
username: root
password: root
# ds1-ds3 类似
rules:
sharding:
tables:
# 订单表:分库分表
t_order:
actual-data-nodes: ds$->{0..3}.t_order_$->{0..15}
# 分库策略(按 user_id)
database-strategy:
standard:
sharding-column: user_id
sharding-algorithm-name: database-inline
# 分表策略(按 order_id)
table-strategy:
standard:
sharding-column: order_id
sharding-algorithm-name: table-inline
# 分布式主键
key-generate-strategy:
column: order_id
key-generator-name: snowflake
# 分片算法
sharding-algorithms:
database-inline:
type: INLINE
props:
algorithm-expression: ds$->{user_id % 4}
table-inline:
type: INLINE
props:
algorithm-expression: t_order_$->{order_id % 16}1.3 自定义分片算法
java
// 自定义复合分片算法
public class DateRangeShardingAlgorithm
implements StandardShardingAlgorithm<LocalDateTime> {
@Override
public String doSharding(Collection<String> availableTargetNames,
PreciseShardingValue<LocalDateTime> shardingValue) {
// 按月份路由:2026-07 → order_202607
LocalDateTime date = shardingValue.getValue();
String suffix = DateTimeFormatter.ofPattern("yyyyMM").format(date);
String target = "t_order_" + suffix;
if (availableTargetNames.contains(target)) {
return target;
}
throw new UnsupportedOperationException("Table " + target + " not exists");
}
@Override
public Collection<String> doSharding(Collection<String> availableTargetNames,
RangeShardingValue<LocalDateTime> shardingValue) {
// 范围查询:跨多个月份
Range<LocalDateTime> range = shardingValue.getValueRange();
LocalDateTime start = range.lowerEndpoint();
LocalDateTime end = range.upperEndpoint();
Collection<String> result = new LinkedList<>();
LocalDateTime current = start;
while (current.isBefore(end) || current.equals(end)) {
String suffix = DateTimeFormatter.ofPattern("yyyyMM").format(current);
String target = "t_order_" + suffix;
if (availableTargetNames.contains(target)) {
result.add(target);
}
current = current.plusMonths(1);
}
return result;
}
}yaml
# 自定义算法注册
sharding-algorithms:
date-range:
type: CLASS_BASED
props:
strategy: STANDARD
algorithmClassName: com.example.DateRangeShardingAlgorithm二、单库分表 + 多库分表混合
2.1 混合策略设计
yaml
rules:
sharding:
tables:
# 订单表:分库 + 分表
t_order:
actual-data-nodes: ds$->{0..3}.t_order_$->{0..15}
# 订单明细:只分库(随主表)
t_order_item:
actual-data-nodes: ds$->{0..3}.t_order_item
database-strategy:
standard:
sharding-column: order_id
sharding-algorithm-name: database-inline
# 用户表:只分表
t_user:
actual-data-nodes: ds0.t_user_$->{0..7}
table-strategy:
standard:
sharding-column: user_id
sharding-algorithm-name: table-inline2.2 广播表与绑定表
yaml
rules:
sharding:
binding-tables: # 绑定表:关联查询避免笛卡尔积
- t_order, t_order_item
- t_product, t_product_detail
broadcast-tables: # 广播表:全库全表一致
- t_dict # 字典表
- t_config # 配置表
- t_region # 地区表三、跨分片查询优化
3.1 问题分析
sql
-- 跨分片分页查询(性能极差)
SELECT * FROM t_order
WHERE status = 'PAID'
ORDER BY create_time DESC
LIMIT 10000, 20;
-- 问题:需要从所有分片取 10020 条,归并后取最后 20 条
-- 分片越多,性能越差3.2 优化方案
java
// 方案 1:游标分页(推荐)
@Mapper
public interface OrderMapper {
@Select("SELECT * FROM t_order " +
"WHERE status = #{status} AND create_time < #{cursor} " +
"ORDER BY create_time DESC LIMIT #{size}")
List<Order> pageByCursor(@Param("status") String status,
@Param("cursor") LocalDateTime cursor,
@Param("size") int size);
// 使用:/api/orders?status=PAID&cursor=2026-07-20T00:00:00&size=20
}
// 方案 2:ES 搜索引擎
@Autowired
private ElasticsearchRestTemplate esTemplate;
public Page<Order> pageByES(String status, int page, int size) {
// ES 集中存储所有订单数据
// 查询 ES 获取 order_id 列表,再根据 order_id 回表查询
NativeQuery query = NativeQuery.builder()
.withQuery(q -> q.term(t -> t.field("status").value(status)))
.withPageable(PageRequest.of(page, size))
.build();
SearchHits<OrderIndex> hits = esTemplate.search(query, OrderIndex.class);
List<Long> ids = hits.stream()
.map(h -> h.getContent().getOrderId())
.collect(Collectors.toList());
// 回表查询(精确路由到具体分片)
return orderMapper.selectBatchIds(ids);
}3.3 全局索引表
java
// 全局索引表:解决买家/卖家双向查询
// t_order_index:order_id + user_id + seller_id
@Table("t_order_index")
public class OrderIndex {
private Long orderId;
private Long userId;
private Long sellerId;
}
// 分片策略:按 user_id + seller_id 双路由
@Repository
public class OrderIndexRepository {
// 买家查订单:按 user_id 路由
public List<Long> findOrderIdsByUserId(Long userId, Pageable pageable) {
// SELECT order_id FROM t_order_index WHERE user_id = ?
}
// 卖家查订单:按 seller_id 路由
public List<Long> findOrderIdsBySellerId(Long sellerId, Pageable pageable) {
// SELECT order_id FROM t_order_index WHERE seller_id = ?
}
// 根据 order_id 回查完整数据
public List<Order> findOrdersByIds(List<Long> orderIds) {
// SELECT * FROM t_order WHERE order_id IN (?)
// ShardingSphere 自动路由到正确分片
}
}四、分布式事务
4.1 配置
yaml
spring:
shardingsphere:
rules:
transaction:
default-type: BASE # 默认使用 BASE 事务
provider-type: Seata # 使用 Seata AT 模式
props:
xa-transaction-manager-type: Atomikos # XA 事务管理器4.2 代码使用
java
@Service
public class OrderServiceImpl {
@Transactional // 本地事务(跨分片不支持强一致)
public void updateStatus(Long orderId) {
orderMapper.updateStatus(orderId, "PAID");
// 如果 orderId 不同分片,本地事务无法保证
}
@ShardingSphereTransactionType(TransactionType.BASE) // 分布式事务
@Transactional
public void createOrderDistributed(Order order) {
// 订单落库
orderMapper.insert(order);
// 扣减库存(跨库操作)
stockMapper.decrease(order.getProductId(), order.getQuantity());
// 积分发放(跨库操作)
pointMapper.add(order.getUserId(), order.getTotalAmount());
}
}五、弹性迁移与扩容
5.1 扩容方案
text
场景:4 库 → 8 库扩容
方案 1:停机迁移(最小改动)
1. 停机维护
2. 导出全部数据
3. 按新分片规则重新入库
4. 修改配置重启
方案 2:双写迁移(不停机)
1. 增加新库(ds4-ds7)
2. 应用双写:同时写入旧库和新库
3. 旧数据通过迁移工具同步到新库
4. 切读:读流量逐渐从旧库切到新库
5. 验证通过后,下线旧库
方案 3:ShardingSphere-Proxy 弹性迁移
使用 Proxy 的内置迁移动能,自动完成数据迁移5.2 分片算法兼容
java
// 兼容旧分片算法:从 4 分片 → 8 分片
// 方案:使用 hash + mod 复合
public class CompatibleShardingAlgorithm implements StandardShardingAlgorithm<Long> {
@Override
public String doSharding(Collection<String> availableTargetNames,
PreciseShardingValue<Long> shardingValue) {
Long value = shardingValue.getValue();
int oldMod = 4;
int newMod = 8;
// 先取 hash 范围
int hash = value.hashCode();
// 旧数据:hash % 4
// 新数据:hash % 8
int target = hash % newMod;
String targetDs = "ds" + target;
if (availableTargetNames.contains(targetDs)) {
return targetDs;
}
// 回退到旧分片
return "ds" + (hash % oldMod);
}
}六、总结
| 知识点 | 说明 |
|---|---|
| 分片策略 | Inline / Standard / Complex / Hint |
| 自定义算法 | CLASS_BASED 实现 StandardShardingAlgorithm |
| 绑定表 | 避免关联查询笛卡尔积 |
| 广播表 | 字典/配置表全库一致 |
| 跨分片分页 | 游标分页 / ES 搜索引擎 / 全局索引 |
| 分布式事务 | XA / Seata BASE 事务 |
| 弹性迁移 | 双写 / Proxy 迁移 / 兼容分片算法 |
参考链接: