Java分库分表案例

wen java案例 3

本文目录导读:

Java分库分表案例

  1. 项目背景
  2. 技术选型
  3. 数据库设计
  4. 核心实现
  5. 测试代码
  6. 运维与监控
  7. 注意事项和最佳实践

这是一个全面的Java分库分表实战案例,包含核心代码和完整实现:

项目背景

假设我们正在开发一个电商系统,订单表数据量巨大(预计每天新增100万条),需要分库分表解决性能问题。

技术选型

  • ShardingSphere-JDBC 5.x(主流分库分表中间件)
  • Spring Boot 2.7
  • MyBatis-Plus
  • MySQL 8.0

数据库设计

分片策略

分库:user_id % 2(2个库)
分表:order_id % 4(每个库4张表)

数据库结构

db_order_0: order_0, order_1, order_2, order_3
db_order_1: order_0, order_1, order_2, order_3

核心实现

Maven依赖(pom.xml)

<dependencies>
    <dependency>
        <groupId>org.apache.shardingsphere</groupId>
        <artifactId>shardingsphere-jdbc-core-spring-boot-starter</artifactId>
        <version>5.3.2</version>
    </dependency>
    <dependency>
        <groupId>com.baomidou</groupId>
        <artifactId>mybatis-plus-boot-starter</artifactId>
        <version>3.5.3</version>
    </dependency>
    <dependency>
        <groupId>mysql</groupId>
        <artifactId>mysql-connector-java</artifactId>
        <version>8.0.33</version>
    </dependency>
</dependencies>

配置文件(application.yml)

spring:
  shardingsphere:
    datasource:
      names: ds0,ds1
      ds0:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://localhost:3306/db_order_0?useSSL=false&serverTimezone=UTC
        username: root
        password: password
      ds1:
        type: com.zaxxer.hikari.HikariDataSource
        driver-class-name: com.mysql.cj.jdbc.Driver
        jdbc-url: jdbc:mysql://localhost:3306/db_order_1?useSSL=false&serverTimezone=UTC
        username: root
        password: password
    rules:
      sharding:
        tables:
          t_order:
            actual-data-nodes: ds$->{0..1}.t_order_$->{0..3}
            table-strategy:
              standard:
                sharding-column: order_id
                sharding-algorithm-name: order_table_inline
            key-generate-strategy:
              column: order_id
              key-generator-name: snowflake
        default-database-strategy:
          standard:
            sharding-column: user_id
            sharding-algorithm-name: database_inline
        sharding-algorithms:
          database_inline:
            type: INLINE
            props:
              algorithm-expression: ds$->{user_id % 2}
          order_table_inline:
            type: INLINE
            props:
              algorithm-expression: t_order_$->{order_id % 4}
        key-generators:
          snowflake:
            type: SNOWFLAKE
            props:
              worker-id: 1
    props:
      sql-show: true  # 显示SQL日志

实体类(Order.java)

@Data
@TableName("t_order")
public class Order {
    @TableId(type = IdType.INPUT)  // 由ShardingSphere生成
    private Long orderId;
    private Long userId;
    private String orderNo;
    private BigDecimal amount;
    private Integer status;
    private LocalDateTime createTime;
    private LocalDateTime updateTime;
}

Mapper接口

@Mapper
public interface OrderMapper extends BaseMapper<Order> {
    /**
     * 自定义查询 - 按用户ID查询订单
     */
    @Select("SELECT * FROM t_order WHERE user_id = #{userId}")
    List<Order> selectByUserId(Long userId);
    /**
     * 分页查询 - 跨分片
     */
    IPage<Order> selectOrderPage(Page<Order> page, @Param("userId") Long userId, 
                                 @Param("status") Integer status);
}

Service层

@Service
@Slf4j
public class OrderServiceImpl implements OrderService {
    @Resource
    private OrderMapper orderMapper;
    @Override
    @Transactional
    public Order createOrder(Order order) {
        // 设置订单号
        order.setOrderNo(generateOrderNo());
        order.setStatus(1);
        order.setCreateTime(LocalDateTime.now());
        order.setUpdateTime(LocalDateTime.now());
        // 插入订单(ShardingSphere自动路由)
        orderMapper.insert(order);
        return order;
    }
    @Override
    public List<Order> queryOrdersByUser(Long userId) {
        // 根据userId路由到对应数据库
        return orderMapper.selectByUserId(userId);
    }
    @Override
    public Page<Order> pageQuery(long pageNum, long pageSize, Long userId) {
        Page<Order> page = new Page<>(pageNum, pageSize);
        return orderMapper.selectOrderPage(page, userId, null);
    }
    /**
     * 生成订单号
     */
    private String generateOrderNo() {
        return "ORD" + System.currentTimeMillis() + RandomUtil.randomNumbers(6);
    }
}

复杂查询XML(OrderMapper.xml)

<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" 
    "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.example.mapper.OrderMapper">
    <!-- 分页查询订单 - 带条件 -->
    <select id="selectOrderPage" resultType="com.example.entity.Order">
        SELECT * FROM t_order
        <where>
            <if test="userId != null">
                AND user_id = #{userId}
            </if>
            <if test="status != null">
                AND status = #{status}
            </if>
        </where>
        ORDER BY create_time DESC
    </select>
    <!-- 统计查询 -->
    <select id="countOrders" resultType="java.lang.Long">
        SELECT COUNT(*) FROM t_order WHERE user_id = #{userId}
    </select>
</mapper>

服务实现类完整版(OrderServiceImpl.java)

@Service
@Slf4j
public class OrderServiceImpl implements OrderService {
    @Resource
    private OrderMapper orderMapper;
    @Resource
    private OrderDetailMapper orderDetailMapper;
    @Override
    @Transactional(rollbackFor = Exception.class)
    public Order createOrder(OrderCreateDTO dto) {
        try {
            // 创建主订单
            Order order = new Order();
            order.setUserId(dto.getUserId());
            order.setAmount(dto.getAmount());
            order.setStatus(OrderStatus.PENDING_PAYMENT.getCode());
            order.setCreateTime(LocalDateTime.now());
            order.setUpdateTime(LocalDateTime.now());
            order.setOrderNo(generateOrderNo());
            orderMapper.insert(order);
            // 创建订单明细
            if (CollectionUtil.isNotEmpty(dto.getItems())) {
                for (OrderItemDTO item : dto.getItems()) {
                    OrderDetail detail = new OrderDetail();
                    detail.setOrderId(order.getOrderId());
                    detail.setProductId(item.getProductId());
                    detail.setQuantity(item.getQuantity());
                    detail.setPrice(item.getPrice());
                    orderDetailMapper.insert(detail);
                }
            }
            log.info("订单创建成功: orderId={}, userId={}", order.getOrderId(), order.getUserId());
            return order;
        } catch (Exception e) {
            log.error("订单创建失败", e);
            throw new BusinessException("订单创建失败");
        }
    }
    @Override
    public Order queryOrder(Long orderId, Long userId) {
        // 关键是传入userId,用于路由到正确的数据库
        if (userId == null) {
            throw new IllegalArgumentException("userId不能为空");
        }
        // 由于分库键是user_id,必须携带userId
        Order order = orderMapper.selectById(orderId);
        if (order == null || !userId.equals(order.getUserId())) {
            throw new BusinessException("订单不存在");
        }
        return order;
    }
    @Override
    public Order updateStatus(Long orderId, Integer status) {
        Order order = orderMapper.selectById(orderId);
        if (order == null) {
            throw new BusinessException("订单不存在");
        }
        order.setStatus(status);
        order.setUpdateTime(LocalDateTime.now());
        orderMapper.updateById(order);
        return order;
    }
}

批量操作示例

@Service
public class BatchOrderService {
    @Resource
    private OrderMapper orderMapper;
    /**
     * 批量插入订单 - 注意:批量操作会并行路由
     */
    @Transactional
    public void batchInsert(List<Order> orders) {
        // 方式1:循环单条插入(推荐,能正确路由)
        for (Order order : orders) {
            orderMapper.insert(order);
        }
        // 方式2:批量插入(需要特殊配置)
        // orderMapper.batchInsert(orders);
    }
    /**
     * 跨分片查询 - 使用Union All
     */
    @Select("SELECT * FROM t_order WHERE create_time BETWEEN #{start} AND #{end}")
    List<Order> queryByTimeRange(LocalDateTime start, LocalDateTime end);
}

分片算法扩展

// 自定义分片算法(可选)
public class CustomShardingAlgorithm implements StandardShardingAlgorithm<Long> {
    @Override
    public String doSharding(Collection<String> availableTargetNames, 
                            PreciseShardingValue<Long> shardingValue) {
        Long value = shardingValue.getValue();
        // 自定义路由逻辑
        if (value % 2 == 0) {
            return "ds0";
        } else {
            return "ds1";
        }
    }
    @Override
    public Collection<String> doSharding(Collection<String> availableTargetNames, 
                                        RangeShardingValue<Long> shardingValue) {
        // 处理范围查询
        Collection<String> result = new ArrayList<>();
        // 实现范围路由逻辑
        return result;
    }
    @Override
    public void init(Properties props) {
        // 初始化配置
    }
    @Override
    public String getType() {
        return "CUSTOM";
    }
}

事务处理

@Configuration
@EnableTransactionManagement
public class TransactionConfig {
    @Bean
    public TransactionTemplate transactionTemplate(
            PlatformTransactionManager transactionManager) {
        return new TransactionTemplate(transactionManager);
    }
}

测试代码

@SpringBootTest
@Slf4j
public class ShardingTest {
    @Resource
    private OrderService orderService;
    @Test
    public void testCreateOrder() {
        OrderCreateDTO dto = new OrderCreateDTO();
        dto.setUserId(1001L);
        dto.setAmount(new BigDecimal("299.00"));
        Order order = orderService.createOrder(dto);
        log.info("创建订单成功: orderId={}", order.getOrderId());
        Assert.notNull(order.getOrderId(), "订单ID不应为空");
    }
    @Test
    public void testQueryOrder() {
        // 使用正确的userId才能路由到正确的库
        Order order = orderService.queryOrder(123456789L, 1001L);
        Assert.notNull(order, "订单不应为空");
    }
    @Test
    public void testQueryOrdersByUser() {
        List<Order> orders = orderService.queryOrdersByUser(1001L);
        Assert.isTrue(orders.size() > 0, "用户应有订单");
    }
    @Test
    public void testPageQuery() {
        Page<Order> page = orderService.pageQuery(1, 10, 1001L);
        Assert.notNull(page.getRecords(), "分页结果不应为空");
    }
}

运维与监控

@Component
public class ShardingMetricCollector {
    @Resource
    private ShardingSphereDataSource dataSource;
    @Scheduled(fixedDelay = 60000)  // 每分钟执行
    public void collectMetrics() {
        // 获取ShardingSphere的统计信息
        Map<String, Object> metrics = dataSource.getRuntimeContext()
            .getStatistics()
            .getDatabaseStatistics();
        // 记录分片执行情况
        log.info("总SQL执行次数: {}", metrics.get("total_sql_count"));
        log.info("分片执行次数: {}", metrics.get("sharding_sql_count"));
    }
}

注意事项和最佳实践

分片键选择

// 必须包含分片键,否则会导致全路由
Order order = new Order();
order.setUserId(userId);  // 必须设置
order.setOrderId(orderId); // 可选

分布式ID生成

// ShardingSphere自带的雪花算法
KeyGenerator keyGenerator = new SnowflakeKeyGenerator();
long id = keyGenerator.generateKey().longValue();

子查询和JOIN支持

<!-- 支持E-R关联查询 -->
<select id="selectOrderWithDetail">
    SELECT o.*, d.product_name, d.quantity
    FROM t_order o
    LEFT JOIN t_order_detail d ON o.order_id = d.order_id
    WHERE o.user_id = #{userId}
</select>

扩容方案

// 推荐使用一致性哈希算法,扩展性更好
PRODUCE: 
- 使用5.x版本的一致哈希分片算法
- 预留分片槽位
- 平滑扩容能力

这个案例涵盖了Java分库分表的主要实现方式,包括配置、代码实现、测试、运维等关键环节,实际项目中还需要根据具体业务场景调整分片策略和配置。

抱歉,评论功能暂时关闭!