MySQL分库分表实战:ShardingSphere-JDBC配置与平滑迁移方案

MySQL在单表数据量超过千万行后查询性能明显下降,B+树层级增加导致磁盘IO放大,索引维护成本急剧上升。分库分表是解决单库容量和性能瓶颈的有效手段,Apache ShardingSphere-JDBC作为轻量级Java框架,以JDBC驱动方式嵌入应用,无需独立部署代理层。本文完整演示ShardingSphere-JDBC的分片配置、数据迁移和运维方案。

分片策略设计与评估

分片策略直接决定数据分布均匀度和查询路由效率。以订单系统为例,按用户ID取模分库、按创建时间分表是常见方案。

// 分片策略评估维度:
// 1. 数据均匀度: 各分片数据量偏差不超过20%
// 2. 查询隔离度: 80%以上查询能路由到单分片
// 3. 扩容成本: 新增分片时数据迁移量
// 4. 跨片查询: 尽量避免跨库JOIN和聚合

// 订单表分片方案:
// 分库: 4个库,按 user_id % 4
// 分表: 每库12张表,按 create_time 月份
// 总分片: 4 × 12 = 48个物理表

ShardingSphere-JDBC Spring Boot配置

# application.yml
spring:
  shardingsphere:
    # 数据源配置
    datasource:
      names: ds0,ds1,ds2,ds3
      ds0:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://10.0.1.10:3306/order_db_0?useUnicode=true&characterEncoding=utf8&serverTimezone=Asia/Shanghai
        username: order_app
        password: ${DB_PASSWORD}
        connection-timeout: 30000
        idle-timeout: 600000
        max-lifetime: 1800000
        maximum-pool-size: 50
      ds1:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://10.0.1.11:3306/order_db_1?useUnicode=true&characterEncoding=utf8&serverTimezone=Asia/Shanghai
        username: order_app
        password: ${DB_PASSWORD}
        maximum-pool-size: 50
      ds2:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://10.0.1.12:3306/order_db_2?useUnicode=true&characterEncoding=utf8&serverTimezone=Asia/Shanghai
        username: order_app
        password: ${DB_PASSWORD}
        maximum-pool-size: 50
      ds3:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://10.0.1.13:3306/order_db_3?useUnicode=true&characterEncoding=utf8&serverTimezone=Asia/Shanghai
        username: order_app
        password: ${DB_PASSWORD}
        maximum-pool-size: 50

    # 分片规则配置
    rules:
      sharding:
        tables:
          t_order:
            # 真实数据节点:4库 × 12表
            actual-data-nodes: ds$->{0..3}.t_order_$->{202601..202612}
            # 分库策略:user_id取模
            database-strategy:
              standard:
                sharding-column: user_id
                sharding-algorithm-name: db-mod
            # 分表策略:按order_date月份
            table-strategy:
              standard:
                sharding-column: order_date
                sharding-algorithm-name: table-month
            # 主键生成策略
            key-generate-strategy:
              column: order_id
              key-generator-name: snowflake
          
          # 订单明细表,与订单表绑定分片
          t_order_item:
            actual-data-nodes: ds$->{0..3}.t_order_item_$->{202601..202612}
            database-strategy:
              standard:
                sharding-column: user_id
                sharding-algorithm-name: db-mod
            table-strategy:
              standard:
                sharding-column: order_date
                sharding-algorithm-name: table-month
            key-generate-strategy:
              column: item_id
              key-generator-name: snowflake

        # 绑定表:同一分片键的表路由到同一数据源
        binding-tables:
          - t_order,t_order_item
        
        # 广播表:全量数据同步到所有库
        broadcast-tables:
          - t_dict_order_status

        # 分片算法定义
        sharding-algorithms:
          db-mod:
            type: MOD
            props:
              sharding-count: 4
          table-month:
            type: INTERVAL
            props:
              datetime-pattern: yyyy-MM-dd
              sharding-suffix-pattern: yyyyMM
              datetime-lower: 2026-01-01
              datetime-upper: 2026-12-31

        # 主键生成器
        key-generators:
          snowflake:
            type: SNOWFLAKE
            props:
              worker-id: 1

    props:
      sql-show: true  # 打印实际路由SQL,生产环境关闭

Java代码中的分片查询

ShardingSphere-JDBC对应用透明,MyBatis或JPA无需修改SQL。但分片键的使用方式直接影响查询效率:

// 包含分片键的查询 - 精确路由到单分片
@Mapper
public interface OrderMapper {

    // 精确路由:user_id + order_date 双分片键
    @Select("SELECT * FROM t_order " +
            "WHERE user_id = #{userId} " +
            "AND order_date = #{orderDate}")
    Order findByUserIdAndDate(@Param("userId") Long userId,
                              @Param("orderDate") LocalDate orderDate);

    // 单分片键路由:仅user_id,需扫描该用户所有月份表
    @Select("SELECT * FROM t_order " +
            "WHERE user_id = #{userId} " +
            "AND create_time BETWEEN #{start} AND #{end}")
    List<Order> findByUserAndTimeRange(@Param("userId") Long userId,
            @Param("start") LocalDateTime start,
            @Param("end") LocalDateTime end);

    // 无分片键查询 - 全分片扫描,性能差
    @Select("SELECT * FROM t_order " +
            "WHERE order_status = #{status} " +
            "AND create_time BETWEEN #{start} AND #{end}")
    List<Order> findByStatusAndTime(@Param("status") Integer status,
            @Param("start") LocalDateTime start,
            @Param("end") LocalDateTime end);

    // 分页查询 - 需要分片合并
    @Select("SELECT * FROM t_order " +
            "WHERE user_id = #{userId} " +
            "ORDER BY create_time DESC " +
            "LIMIT #{offset}, #{size}")
    List<Order> findPageByUser(@Param("userId") Long userId,
            @Param("offset") int offset,
            @Param("size") int size);
}

// 跨分片聚合查询使用Hint强制路由
@Service
public class OrderReportService {

    @Autowired
    private OrderMapper orderMapper;

    public BigDecimal calculateTotalAmount(LocalDateTime start,
            LocalDateTime end) {
        // 使用HintManager指定查询特定分片
        try (HintManager hintManager = HintManager.getInstance()) {
            hintManager.addTableShardingValue("t_order", "202607");
            // 聚合查询在各分片执行后合并结果
            return orderMapper.sumAmount(start, end);
        }
    }
}

数据平滑迁移方案

从单库迁移到分库分表是高风险操作,采用双写+数据同步方案实现平滑迁移:

// 迁移阶段一:双写 + 全量同步
@Service
public class OrderDualWriteService {

    @Autowired
    private OldOrderRepository oldRepo;  // 旧单库
    @Autowired
    private NewOrderMapper newMapper;    // 新分片库
    @Autowired
    private MigrationLogRepository logRepo;

    // 双写入口
    @Transactional
    public void createOrder(Order order) {
        // 1. 写旧库(主)
        oldRepo.save(order);
        
        // 2. 写新库(从)
        try {
            newMapper.insert(order);
            logRepo.save(MigrationLog.success(order.getId()));
        } catch (Exception e) {
            logRepo.save(MigrationLog.failed(order.getId(), e.getMessage()));
            // 异步重试
            retryExecutor.execute(() -> {
                newMapper.insert(order);
            });
        }
    }
}

// 全量数据同步脚本
@Component
public class FullDataSyncJob {

    @Autowired
    private OldOrderRepository oldRepo;
    @Autowired
    private NewOrderMapper newMapper;

    @Scheduled(initialDelay = 60000, fixedDelay = Long.MAX_VALUE)
    public void syncAll() {
        int pageSize = 5000;
        long lastId = 0;
        int total = 0;
        
        while (true) {
            List<Order> batch = oldRepo.findByIdGreaterThan(lastId, pageSize);
            if (batch.isEmpty()) break;
            
            for (Order order : batch) {
                try {
                    // 检查是否已同步
                    if (newMapper.existsById(order.getId())) {
                        continue;
                    }
                    newMapper.insert(order);
                    total++;
                } catch (Exception e) {
                    log.error("同步失败, orderId={}", order.getId(), e);
                }
            }
            
            lastId = batch.get(batch.size() - 1).getId();
            log.info("已同步 {} 条,当前ID: {}", total, lastId);
        }
        log.info("全量同步完成,总计 {} 条", total);
    }
}

// 迁移阶段二:增量同步 + 数据校验
@Component
public class IncrementalSyncJob {

    @Scheduled(fixedDelay = 10000)
    public void syncIncrement() {
        // 基于binlog或更新时间戳同步增量数据
        LocalDateTime lastSyncTime = getLastSyncTime();
        List<Order> changed = oldRepo
            .findByUpdateTimeAfter(lastSyncTime);
        
        for (Order order : changed) {
            Order existing = newMapper.findById(order.getId());
            if (existing == null || 
                !existing.getUpdateTime().equals(order.getUpdateTime())) {
                newMapper.upsert(order);
            }
        }
        
        updateLastSyncTime(LocalDateTime.now());
    }

    // 数据一致性校验
    @Scheduled(cron = "0 0 3 * * ?")
    public void verifyConsistency() {
        int batchSize = 10000;
        long lastId = 0;
        int mismatchCount = 0;
        
        while (true) {
            List<Order> oldBatch = oldRepo
                .findByIdGreaterThan(lastId, batchSize);
            if (oldBatch.isEmpty()) break;
            
            for (Order oldOrder : oldBatch) {
                Order newOrder = newMapper.findById(oldOrder.getId());
                if (newOrder == null || 
                    !dataEquals(oldOrder, newOrder)) {
                    mismatchCount++;
                    log.error("数据不一致, orderId={}, old={}, new={}",
                        oldOrder.getId(), oldOrder, newOrder);
                }
            }
            
            lastId = oldBatch.get(oldBatch.size() - 1).getId();
        }
        
        log.info("校验完成,不一致数量: {}", mismatchCount);
        if (mismatchCount > 0) {
            alertService.notifyDataMismatch(mismatchCount);
        }
    }
}

迁移流程总结:阶段一双写+全量同步 → 阶段二增量同步+数据校验 → 阶段三读写切换到新库 → 阶段四下线旧库。每个阶段都需要预留回滚能力,任何阶段发现问题都可以退回上一阶段。数据校验是迁移过程中最关键的环节,建议在业务低峰期执行全量校验,并持续运行增量校验直到切换完成。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/mysql-fen-ku-fen-biao-shi-zhan-shardingspherejdbc-pei-zhi/

(0)
小编小编
上一篇 9小时前
下一篇 9小时前

相关推荐