datai/docs/implementation/dynamic-datasource-solution.md

23 KiB
Raw Permalink Blame History

延迟加载从库解决方案

需求概述

业务场景

主库(固定)+ 从库(库名可变)
├─ 主库:存储配置表(包含从库库名)
├─ 从库:启动时不加载,运行时根据配置表动态加载
└─ 操作:更新配置表 → 动态加载从库 → 切换写入数据

核心需求

  1. 延迟加载:应用启动时只加载主库,从库按需加载
  2. 动态切换:运行时可以随时切换从库库名
  3. 配置管理:从库配置存储在主库,便于管理
  4. 无侵入性:使用现有 @DataSource 注解,无需修改业务代码

架构设计

核心组件

组件 职责 文件位置
sys_datasource_config 存储从库配置信息 主库表
SysDatasourceConfig 配置实体类 datai-system/domain/
SysDatasourceConfigMapper 配置数据访问层 datai-system/mapper/
DataSourceManager 数据源管理器(扩展) datai-framework/manager/
DynamicDataSourceService 动态数据源业务服务 datai-system/service/
DynamicDataSource 动态数据源路由 datai-framework/datasource/
DynamicDataSourceContextHolder 数据源上下文持有者 datai-framework/datasource/

工作流程

1. 应用启动
   └─ 只加载 MASTER 数据源
   └─ SLAVE 数据源不加载

2. 运行时加载从库
   └─ 读取 sys_datasource_config 表
   └─ 创建 SLAVE 数据源
   └─ 注册到 DynamicDataSource

3. 切换从库库名
   └─ 更新 sys_datasource_config 表
   └─ 移除旧 SLAVE 数据源
   └─ 创建新 SLAVE 数据源
   └─ 注册到 DynamicDataSource

4. 使用从库
   └─ @DataSource(DataSourceType.SLAVE)
   └─ AOP 切面切换数据源
   └─ DynamicDataSourceContextHolder 设置上下文
   └─ DynamicDataSource 路由到 SLAVE

实现步骤

步骤 1创建数据源配置表

文件位置: datai-system/src/main/resources/sql/sys_datasource_config.sql

-- 数据源配置表
-- 用于存储从库的配置信息,支持动态切换从库库名
CREATE TABLE `sys_datasource_config` (
  `id` BIGINT(20) NOT NULL AUTO_INCREMENT COMMENT '主键ID',
  `ds_name` VARCHAR(50) NOT NULL COMMENT '数据源名称SLAVE',
  `db_name` VARCHAR(100) NOT NULL COMMENT '数据库名称',
  `db_host` VARCHAR(200) NOT NULL COMMENT '数据库主机',
  `db_port` INT(11) DEFAULT '3306' COMMENT '数据库端口',
  `username` VARCHAR(100) NOT NULL COMMENT '用户名',
  `password` VARCHAR(200) NOT NULL COMMENT '密码',
  `db_type` VARCHAR(20) DEFAULT 'mysql' COMMENT '数据库类型mysql/postgresql/oracle等',
  `status` CHAR(1) DEFAULT '0' COMMENT '状态0正常 1停用',
  `remark` VARCHAR(500) DEFAULT NULL COMMENT '备注',
  `create_by` VARCHAR(64) DEFAULT '' COMMENT '创建者',
  `create_time` DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
  `update_by` VARCHAR(64) DEFAULT '' COMMENT '更新者',
  `update_time` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_ds_name` (`ds_name`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='数据源配置表';

-- 初始化从库配置
INSERT INTO `sys_datasource_config` (`ds_name`, `db_name`, `db_host`, `db_port`, `username`, `password`, `db_type`, `status`, `remark`)
VALUES ('SLAVE', 'datai_slave', '117.72.193.251', 3306, 'root', '2026.DataI', 'mysql', '0', '从库配置');

步骤 2修改配置文件

文件位置: datai-admin/src/main/resources/application-druid.yml

spring:
  datasource:
    type: com.alibaba.druid.pool.DruidDataSource
    driverClassName: com.mysql.cj.jdbc.Driver
    dynamic:
      primary: MASTER
      xa: false  # 暂不启用分布式事务
      datasource:
        MASTER:  # 只配置主库
          url: jdbc:mysql://117.72.193.251/datai?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&useSSL=true&serverTimezone=GMT%2B8&nullCatalogMeansCurrent=true
          username: root
          password: 2026.DataI
        # SLAVE:  # 移除从库配置,运行时动态加载
    druid:
      initialSize: 5
      minIdle: 10
      maxActive: 20
      maxWait: 60000
      connectTimeout: 30000
      socketTimeout: 60000
      timeBetweenEvictionRunsMillis: 60000
      minEvictableIdleTimeMillis: 300000
      maxEvictableIdleTimeMillis: 900000
      validationQuery: SELECT 1
      testWhileIdle: true
      testOnBorrow: false
      testOnReturn: false
      filters: stat,wall,slf4j,log4j2
      connectionProperties: druid.stat.mergeSql=true;druid.stat.slowSqlMillis=5000

步骤 3创建实体类

文件位置: datai-system/src/main/java/com/datai/system/domain/SysDatasourceConfig.java

package com.datai.system.domain;

import com.datai.common.core.domain.BaseEntity;

public class SysDatasourceConfig extends BaseEntity {
    private static final long serialVersionUID = 1L;

    private Long id;
    private String dsName;
    private String dbName;
    private String dbHost;
    private Integer dbPort;
    private String username;
    private String password;
    private String dbType;
    private String status;

    public Long getId() {
        return id;
    }

    public void setId(Long id) {
        this.id = id;
    }

    public String getDsName() {
        return dsName;
    }

    public void setDsName(String dsName) {
        this.dsName = dsName;
    }

    public String getDbName() {
        return dbName;
    }

    public void setDbName(String dbName) {
        this.dbName = dbName;
    }

    public String getDbHost() {
        return dbHost;
    }

    public void setDbHost(String dbHost) {
        this.dbHost = dbHost;
    }

    public Integer getDbPort() {
        return dbPort;
    }

    public void setDbPort(Integer dbPort) {
        this.dbPort = dbPort;
    }

    public String getUsername() {
        return username;
    }

    public void setUsername(String username) {
        this.username = username;
    }

    public String getPassword() {
        return password;
    }

    public void setPassword(String password) {
        this.password = password;
    }

    public String getDbType() {
        return dbType;
    }

    public void setDbType(String dbType) {
        this.dbType = dbType;
    }

    public String getStatus() {
        return status;
    }

    public void setStatus(String status) {
        this.status = status;
    }
}

步骤 4创建 Mapper

文件位置: datai-system/src/main/java/com/datai/system/mapper/SysDatasourceConfigMapper.java

package com.datai.system.mapper;

import com.datai.system.domain.SysDatasourceConfig;
import org.apache.ibatis.annotations.Mapper;

@Mapper
public interface SysDatasourceConfigMapper {

    SysDatasourceConfig selectByDsName(String dsName);

    int updateByDsName(SysDatasourceConfig config);

    int insert(SysDatasourceConfig config);
}

文件位置: datai-system/src/main/resources/mapper/system/SysDatasourceConfigMapper.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.datai.system.mapper.SysDatasourceConfigMapper">

    <resultMap type="SysDatasourceConfig" id="SysDatasourceConfigResult">
        <id     property="id"           column="id"           />
        <result property="dsName"       column="ds_name"      />
        <result property="dbName"       column="db_name"      />
        <result property="dbHost"       column="db_host"      />
        <result property="dbPort"       column="db_port"      />
        <result property="username"     column="username"     />
        <result property="password"     column="password"     />
        <result property="dbType"       column="db_type"      />
        <result property="status"       column="status"       />
        <result property="remark"       column="remark"       />
        <result property="createBy"     column="create_by"    />
        <result property="createTime"   column="create_time"  />
        <result property="updateBy"     column="update_by"    />
        <result property="updateTime"   column="update_time"  />
    </resultMap>

    <select id="selectByDsName" parameterType="String" resultMap="SysDatasourceConfigResult">
        SELECT * FROM sys_datasource_config WHERE ds_name = #{dsName}
    </select>

    <update id="updateByDsName" parameterType="SysDatasourceConfig">
        UPDATE sys_datasource_config
        <set>
            <if test="dbName != null and dbName != ''">db_name = #{dbName},</if>
            <if test="dbHost != null and dbHost != ''">db_host = #{dbHost},</if>
            <if test="dbPort != null">db_port = #{dbPort},</if>
            <if test="username != null and username != ''">username = #{username},</if>
            <if test="password != null and password != ''">password = #{password},</if>
            <if test="dbType != null and dbType != ''">db_type = #{dbType},</if>
            <if test="status != null and status != ''">status = #{status},</if>
            <if test="remark != null">remark = #{remark},</if>
            update_by = #{updateBy},
            update_time = sysdate()
        </set>
        WHERE ds_name = #{dsName}
    </update>

    <insert id="insert" parameterType="SysDatasourceConfig" useGeneratedKeys="true" keyProperty="id">
        INSERT INTO sys_datasource_config (
            ds_name, db_name, db_host, db_port, username, password, db_type, status, remark, create_by, create_time
        ) VALUES (
            #{dsName}, #{dbName}, #{dbHost}, #{dbPort}, #{username}, #{password}, #{dbType}, #{status}, #{remark}, #{createBy}, sysdate()
        )
    </insert>

</mapper>

步骤 5扩展 DataSourceManager

文件位置: datai-framework/src/main/java/com/datai/framework/manager/DataSourceManager.java

在现有代码中添加以下方法:

public class DataSourceManager implements InitializingBean {
    
    // ... 现有代码 ...

    public synchronized void addDynamicDataSource(String name, String url, String username, String password) {
        if (targetDataSources.containsKey(name)) {
            logger.warn("数据源 {} 已存在,先移除旧数据源", name);
            removeDataSource(name);
        }
        
        Properties prop = new Properties();
        prop.setProperty("druid.url", url);
        prop.setProperty("druid.username", username);
        prop.setProperty("druid.password", password);
        
        DruidProperties druidProperties = SpringUtils.getBean(DruidProperties.class);
        prop.setProperty("druid.initialSize", String.valueOf(druidProperties.getInitialSize()));
        prop.setProperty("druid.minIdle", String.valueOf(druidProperties.getMinIdle()));
        prop.setProperty("druid.maxActive", String.valueOf(druidProperties.getMaxActive()));
        prop.setProperty("druid.maxWait", String.valueOf(druidProperties.getMaxWait()));
        prop.setProperty("druid.validationQuery", druidProperties.getValidationQuery());
        prop.setProperty("druid.testWhileIdle", String.valueOf(druidProperties.isTestWhileIdle()));
        prop.setProperty("druid.testOnBorrow", String.valueOf(druidProperties.isTestOnBorrow()));
        prop.setProperty("druid.testOnReturn", String.valueOf(druidProperties.isTestOnReturn()));
        prop.setProperty("druid.filters", druidProperties.getFilters());
        prop.setProperty("druid.connectionProperties", druidProperties.getConnectionProperties());
        prop.setProperty("druid.timeBetweenEvictionRunsMillis",
                String.valueOf(druidProperties.getTimeBetweenEvictionRunsMillis()));
        prop.setProperty("druid.minEvictableIdleTimeMillis",
                String.valueOf(druidProperties.getMinEvictableIdleTimeMillis()));
        prop.setProperty("druid.maxEvictableIdleTimeMillis",
                String.valueOf(druidProperties.getMaxEvictableIdleTimeMillis()));
        
        DataSource dataSource = createDataSource(name, prop);
        putDataSource(name, dataSource);
        
        logger.info("动态数据源 {} 添加成功", name);
    }
    
    public synchronized void removeDataSource(String name) {
        DataSource dataSource = targetDataSources.remove(name);
        if (dataSource != null) {
            dataSourceKeyIndex.remove(dataSource);
            dsDatabaseId.remove(name);
            if (dataSource instanceof DruidDataSource) {
                try {
                    ((DruidDataSource) dataSource).close();
                } catch (Exception e) {
                    logger.error("关闭数据源 {} 失败", name, e);
                }
            }
            logger.info("数据源 {} 已移除", name);
        }
    }
    
    public synchronized void updateDataSource(String name, String url, String username, String password) {
        removeDataSource(name);
        addDynamicDataSource(name, url, username, password);
    }
}

步骤 6创建动态数据源服务

文件位置: datai-system/src/main/java/com/datai/system/service/DynamicDataSourceService.java

package com.datai.system.service;

import com.datai.common.annotation.DataSource;
import com.datai.common.enums.DataSourceType;
import com.datai.framework.manager.DataSourceManager;
import com.datai.system.domain.SysDatasourceConfig;
import com.datai.system.mapper.SysDatasourceConfigMapper;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

@Service
public class DynamicDataSourceService {
    
    @Autowired
    private DataSourceManager dataSourceManager;
    
    @Autowired
    private SysDatasourceConfigMapper configMapper;
    
    @DataSource(DataSourceType.MASTER)
    public void loadSlaveDataSource() {
        SysDatasourceConfig config = configMapper.selectByDsName("SLAVE");
        if (config == null) {
            throw new RuntimeException("从库配置不存在");
        }
        
        if (!"0".equals(config.getStatus())) {
            throw new RuntimeException("从库配置已停用");
        }
        
        String url = buildJdbcUrl(config);
        dataSourceManager.addDynamicDataSource(
            config.getDsName(),
            url,
            config.getUsername(),
            config.getPassword()
        );
    }
    
    @DataSource(DataSourceType.MASTER)
    public void updateSlaveConfig(String newDbName) {
        SysDatasourceConfig config = configMapper.selectByDsName("SLAVE");
        if (config == null) {
            throw new RuntimeException("从库配置不存在");
        }
        
        config.setDbName(newDbName);
        configMapper.updateByDsName(config);
        
        String url = buildJdbcUrl(config);
        dataSourceManager.updateDataSource(
            config.getDsName(),
            url,
            config.getUsername(),
            config.getPassword()
        );
    }
    
    @DataSource(DataSourceType.MASTER)
    public SysDatasourceConfig getSlaveConfig() {
        return configMapper.selectByDsName("SLAVE");
    }
    
    private String buildJdbcUrl(SysDatasourceConfig config) {
        String urlPattern;
        switch (config.getDbType().toLowerCase()) {
            case "mysql":
                urlPattern = "jdbc:mysql://%s:%d/%s?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&useSSL=true&serverTimezone=GMT%%2B8";
                break;
            case "postgresql":
                urlPattern = "jdbc:postgresql://%s:%d/%s";
                break;
            case "oracle":
                urlPattern = "jdbc:oracle:thin:@%s:%d:%s";
                break;
            default:
                throw new RuntimeException("不支持的数据库类型:" + config.getDbType());
        }
        return String.format(urlPattern, config.getDbHost(), config.getDbPort(), config.getDbName());
    }
}

步骤 7创建控制器可选

文件位置: datai-admin/src/main/java/com/datai/web/controller/system/DatasourceController.java

package com.datai.web.controller.system;

import com.datai.common.annotation.DataSource;
import com.datai.common.core.controller.BaseController;
import com.datai.common.core.domain.AjaxResult;
import com.datai.common.enums.DataSourceType;
import com.datai.system.domain.SysDatasourceConfig;
import com.datai.system.service.DynamicDataSourceService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;

@RestController
@RequestMapping("/system/datasource")
public class DatasourceController extends BaseController {
    
    @Autowired
    private DynamicDataSourceService dynamicDataSourceService;
    
    @PostMapping("/loadSlave")
    public AjaxResult loadSlave() {
        dynamicDataSourceService.loadSlaveDataSource();
        return AjaxResult.success("从库加载成功");
    }
    
    @PostMapping("/switchSlave/{dbName}")
    public AjaxResult switchSlave(@PathVariable String dbName) {
        dynamicDataSourceService.updateSlaveConfig(dbName);
        return AjaxResult.success("从库切换成功");
    }
    
    @GetMapping("/getSlaveConfig")
    @DataSource(DataSourceType.MASTER)
    public AjaxResult getSlaveConfig() {
        SysDatasourceConfig config = dynamicDataSourceService.getSlaveConfig();
        return AjaxResult.success(config);
    }
}

使用示例

示例 1启动时加载从库

@Component
public class DataSourceInitializer implements ApplicationRunner {
    
    @Autowired
    private DynamicDataSourceService dynamicDataSourceService;
    
    @Override
    public void run(ApplicationArguments args) throws Exception {
        try {
            dynamicDataSourceService.loadSlaveDataSource();
            logger.info("从库数据源加载成功");
        } catch (Exception e) {
            logger.warn("从库数据源加载失败,将延迟加载", e);
        }
    }
}

示例 2切换从库库名

@Service
public class OrderService {
    
    @Autowired
    private DynamicDataSourceService dynamicDataSourceService;
    
    @Autowired
    private OrderMapper orderMapper;
    
    @DataSource(DataSourceType.MASTER)
    public void processOrderToSlave(Long orderId, String slaveDbName) {
        Order order = orderMapper.selectById(orderId);
        
        dynamicDataSourceService.updateSlaveConfig(slaveDbName);
        
        writeToSlave(order);
    }
    
    @DataSource(DataSourceType.SLAVE)
    private void writeToSlave(Order order) {
        orderMapper.insert(order);
    }
}

示例 3使用从库写入数据

@Service
public class DataService {
    
    @Autowired
    private DataMapper dataMapper;
    
    @DataSource(DataSourceType.SLAVE)
    public void writeToSlave(Data data) {
        dataMapper.insert(data);
    }
    
    @DataSource(DataSourceType.MASTER)
    public void writeToMaster(Data data) {
        dataMapper.insert(data);
    }
}

注意事项

⚠️ 事务问题

错误做法:

@Transactional  // ❌ 不要在事务中切换数据源
public void switchAndWrite() {
    dynamicDataSourceService.updateSlaveConfig("new_db");
    writeToSlave(data);
}

正确做法:

public void switchAndWrite() {
    dynamicDataSourceService.updateSlaveConfig("new_db");
    writeToSlave(data);  // 另一个方法,独立事务
}

⚠️ 连接池管理

  • 切换数据源时会关闭旧连接池
  • 可能有短暂的服务不可用(毫秒级)
  • 建议在低峰期切换

⚠️ 并发问题

已通过 synchronized 关键字保证线程安全:

public synchronized void updateDataSource(String name, String url, String username, String password) {
    removeDataSource(name);
    addDynamicDataSource(name, url, username, password);
}

⚠️ 数据库类型支持

当前支持:

  • MySQL
  • PostgreSQL
  • Oracle

如需支持其他数据库,在 buildJdbcUrl() 方法中添加对应的 URL 模板。


方案优势

特性 说明
延迟加载 启动时只加载主库,从库按需加载
动态切换 运行时可以随时切换从库库名
配置管理 从库配置存储在主库,便于管理
线程安全 使用 synchronized 和 ThreadLocal 保证线程隔离
无侵入 使用现有 @DataSource 注解,无需修改业务代码
易扩展 支持多种数据库类型

测试验证

测试步骤

  1. 执行 SQL 脚本:创建 sys_datasource_config
  2. 修改配置文件:移除 SLAVE 配置
  3. 启动应用:验证只加载 MASTER 数据源
  4. 调用加载接口POST /system/datasource/loadSlave
  5. 验证从库可用:使用 @DataSource(DataSourceType.SLAVE) 写入数据
  6. 切换从库库名POST /system/datasource/switchSlave/new_db_name
  7. 验证新从库:确认数据写入到新数据库

验证命令

# 1. 加载从库
curl -X POST http://localhost:8080/system/datasource/loadSlave

# 2. 查看从库配置
curl http://localhost:8080/system/datasource/getSlaveConfig

# 3. 切换从库
curl -X POST http://localhost:8080/system/datasource/switchSlave/datai_slave_new

相关文件

文件 路径
配置文件 application-druid.yml
数据源管理器 DataSourceManager.java
动态数据源 DynamicDataSource.java
上下文持有者 DynamicDataSourceContextHolder.java
数据源注解 DataSource.java
数据源类型 DataSourceType.java
AOP 切面 DataSourceAspect.java

总结

本方案实现了延迟加载从库的功能,通过在主库存储从库配置,运行时动态创建和切换数据源,满足了业务需求。方案具有以下特点:

  1. 灵活性高:支持运行时动态切换从库库名
  2. 侵入性低:使用现有注解机制,无需修改业务代码
  3. 可维护性强:配置集中管理,便于运维
  4. 线程安全:通过同步机制保证并发安全

该方案适用于需要动态切换数据源的业务场景,如多租户数据隔离、数据归档、报表库切换等。