RuoYi-Cloud 双写功能实现总结(MySQL ↔ OceanBase)

完整双写代码(17 Java + 1 SQL)位于 ruoyi-modules/ruoyi-system/src/main/java/com/ruoyi/system/dualwrite/


一、包结构总览

com.ruoyi.system.dualwrite/
├── config/                              # 配置类(5个)
│   ├── DualWriteProperties.java             全局配置 @RefreshScope
│   ├── DualWriteOceanbaseProperties.java     OB 连接属性
│   ├── DualWriteDataSourceConfig.java        数据源定义 @Primary + @Bean("oceanbaseDs")
│   ├── DualWriteAutoConfiguration.java       条件装配入口
│   └── AsyncConfig.java                     异步线程池
├── executor/                            # 执行器(3个)
│   ├── DualWriteExecutor.java                @Async 异步/同步执行
│   ├── DualWriteTask.java                    任务载体
│   └── SnowflakeIdGenerator.java             雪花算法 ID
├── interceptor/                         # 拦截器(2个)
│   ├── DualWriteInterceptor.java             核心 @Intercepts
│   └── DualWriteDdlListener.java             DDL 同步
├── router/                              # 路由与健康检查(3个)
│   ├── DualWriteRouter.java                  路由决策 volatile
│   ├── DataSourceHealthIndicator.java        Actuator + @Scheduled 检测
│   └── DataSourceRole.java                   MASTER / SLAVE 枚举
├── log/                                 # 日志(2个)
│   ├── entity/DualWriteLog.java              @Table 实体
│   └── service/DualWriteLogService.java       JdbcTemplate 实现
└── job/
    └── DualWriteLogCompensationJob.java      补偿定时 @Scheduled

二、源码

2.1 config/DualWriteProperties.java — 双写全局配置

package com.ruoyi.system.dualwrite.config;

import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.cloud.context.config.annotation.RefreshScope;

/**
 * 双写全局配置属性类
 *
 * <p>该类的所有属性均以 <b>dual-write.*</b> 为前缀配置在 bootstrap.yml 中,
 * 支持通过 Nacos 配置中心动态刷新 (@RefreshScope),运行时修改无需重启应用。</p>
 *
 * <p>核心配置项包括:</p>
 * <ul>
 *   <li><b>enabled</b> — 总开关,关闭后拦截器不生效,MyBatis 行为与原来完全一致</li>
 *   <li><b>active</b> — 主库角色(MASTER=MySQL / SLAVE=OceanBase),运行时可动态切换</li>
 *   <li><b>syncMode</b> — 同步模式(ASYNC=异步 / SYNC=同步 / MANUAL=手动)</li>
 *   <li><b>autoFailover</b> — 自动故障切换开关</li>
 * </ul>
 *
 * @author ruoyi
 */
@RefreshScope
@ConfigurationProperties(prefix = "dual-write")
public class DualWriteProperties {

    /**
     * 双写功能总开关
     * <p>设置为 true 时启用双写功能,MyBatis 拦截器会捕获所有 INSERT/UPDATE/DELETE 并同步到 OceanBase。
     * 设置为 false 时拦截器直接放行,不产生任何额外开销。</p>
     * <p><b>默认值:</b>false(关闭)</p>
     */
    private boolean enabled = false;

    /**
     * 当前主库角色
     * <p>MASTER 表示 MySQL 为主库(默认),SLAVE 表示 OceanBase 为主库。
     * 故障切换时由 DataSourceHealthIndicator 自动修改此值,
     * 也可以通过手动配置强制切换。</p>
     * <p><b>可选值:</b>MASTER / SLAVE</p>
     * <p><b>默认值:</b>MASTER</p>
     */
    private String active = "MASTER";

    /**
     * 双写同步模式
     * <ul>
     *   <li><b>ASYNC</b>(异步,默认)— 主库执行成功后,将 SQL 提交到独立线程池写入 OceanBase,
     *       不阻塞主库响应。适合对一致性要求不高、对性能敏感的场景。</li>
     *   <li><b>SYNC</b>(同步)— 主库执行成功后,在同一线程中等待 OceanBase 写入完成,
     *       若 OB 写入失败则整个事务回滚。适合对数据一致性要求高的场景。</li>
     *   <li><b>MANUAL</b>(手动)— 仅记录日志到 dual_write_log 表,不自动执行同步,
     *       由补偿定时任务或人工触发。适合灰度验证阶段。</li>
     * </ul>
     * <p><b>默认值:</b>ASYNC</p>
     */
    private String syncMode = "ASYNC";

    /**
     * 是否启用自动故障切换
     * <p>当 OceanBase 健康检查连续失败达到阈值(默认 3 次)时,
     * 是否自动将主库从 MySQL 切换到 OceanBase。</p>
     * <p>关闭时即使 OceanBase 不可用也不会自动切换,需人工介入。</p>
     * <p><b>默认值:</b>true(启用)</p>
     */
    private boolean autoFailover = true;

    /**
     * 雪花算法配置
     * <p>用于生成全局唯一 ID,替代数据库自增主键,避免双写场景下 MySQL 和 OceanBase 产生主键冲突。</p>
     */
    private SnowflakeProperties snowflake = new SnowflakeProperties();

    /**
     * 是否同步 DDL 语句到 OceanBase
     * <p>开启后,MySQL 上执行的 CREATE TABLE / ALTER TABLE / DROP TABLE 等 DDL
     * 会自动在 OceanBase 上执行同样的操作,保持两库表结构一致。</p>
     * <p><b>默认值:</b>true(启用)</p>
     */
    private boolean ddlSync = true;

    /**
     * 双写线程池配置
     * <p>仅在 syncMode=ASYNC 时生效,控制异步写入 OceanBase 的线程池参数。</p>
     */
    private ThreadPoolProperties threadPool = new ThreadPoolProperties();

    /**
     * 失败重试配置
     * <p>控制双写失败时的重试次数和间隔时间。</p>
     */
    private RetryProperties retry = new RetryProperties();

    /**
     * 雪花算法参数
     * <p>每个微服务实例应配置不同的 workerId,避免 ID 冲突。</p>
     */
    public static class SnowflakeProperties {
        /**
         * 工作节点 ID(取值范围 0~31)
         * <p>每个微服务实例应分配不同的 workerId,确保生成的全局 ID 唯一不重复。
         * 如果部署了多个 system 服务实例,需要手动分配不同的 workerId。</p>
         */
        private long workerId = 1;

        /**
         * 数据中心 ID(取值范围 0~31)
         * <p>在多数据中心部署时使用,单机房部署建议保持默认值 0。</p>
         */
        private long datacenterId = 0;

        public long getWorkerId() { return workerId; }
        public void setWorkerId(long workerId) { this.workerId = workerId; }
        public long getDatacenterId() { return datacenterId; }
        public void setDatacenterId(long datacenterId) { this.datacenterId = datacenterId; }
    }

    /**
     * 异步双写线程池参数
     * <p>对应 Spring ThreadPoolTaskExecutor 的配置,仅在异步模式下使用。</p>
     */
    public static class ThreadPoolProperties {
        /** 核心线程数(默认 4),即使空闲也保持存活的线程数量 */
        private int corePoolSize = 4;
        /** 最大线程数(默认 8),队列满时最多能创建的线程数 */
        private int maxPoolSize = 8;
        /** 任务队列容量(默认 200),核心线程满时任务排队的最大长度 */
        private int queueCapacity = 200;

        public int getCorePoolSize() { return corePoolSize; }
        public void setCorePoolSize(int corePoolSize) { this.corePoolSize = corePoolSize; }
        public int getMaxPoolSize() { return maxPoolSize; }
        public void setMaxPoolSize(int maxPoolSize) { this.maxPoolSize = maxPoolSize; }
        public int getQueueCapacity() { return queueCapacity; }
        public void setQueueCapacity(int queueCapacity) { this.queueCapacity = queueCapacity; }
    }

    /**
     * 失败重试参数
     * <p>当双写 OceanBase 失败时,自动重试的配置。</p>
     */
    public static class RetryProperties {
        /** 最大重试次数(默认 3),超过后标记为 FAILED 不再重试 */
        private int maxAttempts = 3;
        /** 重试间隔(默认 1000 毫秒),每次重试之间等待的时间 */
        private long backoffDelay = 1000;

        public int getMaxAttempts() { return maxAttempts; }
        public void setMaxAttempts(int maxAttempts) { this.maxAttempts = maxAttempts; }
        public long getBackoffDelay() { return backoffDelay; }
        public void setBackoffDelay(long backoffDelay) { this.backoffDelay = backoffDelay; }
    }

    public boolean isEnabled() { return enabled; }
    public void setEnabled(boolean enabled) { this.enabled = enabled; }
    public String getActive() { return active; }
    public void setActive(String active) { this.active = active; }
    public String getSyncMode() { return syncMode; }
    public void setSyncMode(String syncMode) { this.syncMode = syncMode; }
    public boolean isAutoFailover() { return autoFailover; }
    public void setAutoFailover(boolean autoFailover) { this.autoFailover = autoFailover; }
    public SnowflakeProperties getSnowflake() { return snowflake; }
    public void setSnowflake(SnowflakeProperties snowflake) { this.snowflake = snowflake; }
    public boolean isDdlSync() { return ddlSync; }
    public void setDdlSync(boolean ddlSync) { this.ddlSync = ddlSync; }
    public ThreadPoolProperties getThreadPool() { return threadPool; }
    public void setThreadPool(ThreadPoolProperties threadPool) { this.threadPool = threadPool; }
    public RetryProperties getRetry() { return retry; }
    public void setRetry(RetryProperties retry) { this.retry = retry; }
}

对应 Nacos 配置:

dual-write:
  enabled: true
  active: MASTER
  sync-mode: ASYNC
  auto-failover: true

2.2 config/DualWriteOceanbaseProperties.java — OB 连接属性

package com.ruoyi.system.dualwrite.config;

import org.springframework.boot.context.properties.ConfigurationProperties;

/**
 * OceanBase 数据源连接属性
 *
 * <p>该类的属性以 <b>oceanbase.datasource.*</b> 为前缀配置在 Nacos 中,
 * 用于创建一个独立的 HikariCP 连接池来连接 OceanBase 数据库。</p>
 *
 * <p>注意:这个数据源与 MyBatis-Flex 管理的主数据源(MySQL)完全隔离,
 * 不会被 MyBatis-Flex 的自动配置影响,也不会出现在 MyBatis-Flex 的
 * 数据源列表中。它是一个纯粹的 JDBC 连接池,仅供双写执行器使用。</p>
 *
 * @author ruoyi
 */
@ConfigurationProperties(prefix = "oceanbase.datasource")
public class DualWriteOceanbaseProperties {

    /**
     * OceanBase JDBC 连接地址
     * <p>OceanBase 兼容 MySQL 协议,因此使用 MySQL JDBC 驱动。
     * 默认连接到 192.168.2.32:2881 的 ry 数据库。
     * useSSL=false 因为 OB 社区版默认不开启 SSL。
     * allowPublicKeyRetrieval=true 允许客户端自动获取公钥进行密码加密。</p>
     */
    private String url = "jdbc:mysql://192.168.2.32:2881/ry?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Shanghai";

    /**
     * OceanBase 登录用户名
     * <p>OceanBase 使用 tenant@user#cluster 格式的用户名。
     * root@sys 表示使用 sys 租户的 root 用户。
     * 如果后续创建了业务租户,应使用对应租户的用户名。</p>
     */
    private String username = "root@sys";

    /**
     * OceanBase 登录密码
     */
    private String password = "123";

    /**
     * JDBC 驱动类名
     * <p>OceanBase 兼容 MySQL 协议,使用 com.mysql.cj.jdbc.Driver。</p>
     */
    private String driverClassName = "com.mysql.cj.jdbc.Driver";

    /**
     * HikariCP 连接池专属配置
     */
    private HikariConfig hikari = new HikariConfig();

    /**
     * HikariCP 连接池参数
     * <p>对 OceanBase 的连接池参数可以独立于主库进行调优。
     * 由于双写是异步的,连接池可以设置得比主库小一些。</p>
     */
    public static class HikariConfig {
        /** 最大连接数(默认 10),OB 不是主库,不需要太多连接 */
        private int maximumPoolSize = 10;
        /** 最小空闲连接数(默认 2) */
        private int minimumIdle = 2;
        /** 连接超时时间(默认 30 秒),获取连接等待的最长时间 */
        private long connectionTimeout = 30000;
        /** 空闲超时时间(默认 10 分钟),空闲连接超过此时间会被回收 */
        private long idleTimeout = 600000;
        /** 连接最大存活时间(默认 30 分钟),防止长时间使用的连接被网络设备断开 */
        private long maxLifetime = 1800000;

        public int getMaximumPoolSize() { return maximumPoolSize; }
        public void setMaximumPoolSize(int maximumPoolSize) { this.maximumPoolSize = maximumPoolSize; }
        public int getMinimumIdle() { return minimumIdle; }
        public void setMinimumIdle(int minimumIdle) { this.minimumIdle = minimumIdle; }
        public long getConnectionTimeout() { return connectionTimeout; }
        public void setConnectionTimeout(long connectionTimeout) { this.connectionTimeout = connectionTimeout; }
        public long getIdleTimeout() { return idleTimeout; }
        public void setIdleTimeout(long idleTimeout) { this.idleTimeout = idleTimeout; }
        public long getMaxLifetime() { return maxLifetime; }
        public void setMaxLifetime(long maxLifetime) { this.maxLifetime = maxLifetime; }
    }

    public String getUrl() { return url; }
    public void setUrl(String url) { this.url = url; }
    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 getDriverClassName() { return driverClassName; }
    public void setDriverClassName(String driverClassName) { this.driverClassName = driverClassName; }
    public HikariConfig getHikari() { return hikari; }
    public void setHikari(HikariConfig hikari) { this.hikari = hikari; }
}

对应 Nacos 配置:

oceanbase:
  datasource:
    url: jdbc:mysql://192.168.2.32:2881/ry?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Shanghai
    username: root@sys
    password: '123'
    driver-class-name: com.mysql.cj.jdbc.Driver
    hikari:
      maximum-pool-size: 10
      minimum-idle: 2

2.3 config/DualWriteDataSourceConfig.java — 数据源定义

package com.ruoyi.system.dualwrite.config;

import com.zaxxer.hikari.HikariConfig;
import com.zaxxer.hikari.HikariDataSource;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.boot.jdbc.DataSourceBuilder;
import org.springframework.boot.jdbc.init.DataSourceScriptDatabaseInitializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Primary;

import javax.sql.DataSource;

/**
 * OceanBase 数据源 Bean 定义
 *
 * <p>该配置类仅在 <b>dual-write.enabled=true</b> 时生效,
 * 创建一个名为 "oceanbaseDs" 的独立 HikariCP 连接池。</p>
 *
 * <p>为什么需要一个独立的数据源而不是使用 MyBatis-Flex 的多数据源功能?</p>
 * <ul>
 *   <li>MyBatis-Flex 的多数据源是基于每条 SQL 路由到不同的数据源,
 *       而双写的需求是同一条 SQL 同时写入两个数据源。</li>
 *   <li>独立数据源不会被 MyBatis-Flex 的拦截器和缓存影响,
 *       保持纯粹的 JDBC 操作。</li>
 *   <li>独立连接池可以独立配置参数(如连接数量),不影响主库性能。</li>
 * </ul>
 *
 * @author ruoyi
 */
@Configuration
@EnableConfigurationProperties(DualWriteOceanbaseProperties.class)
@ConditionalOnProperty(prefix = "dual-write", name = "enabled", havingValue = "true", matchIfMissing = false)
public class DualWriteDataSourceConfig {

    /**
     * 创建 MySQL 主数据源(@Primary)
     *
     * <p>从 spring.datasource.* 配置创建,保证 MyBatis-Flex 和 JdbcTemplate
     * 使用 MySQL 作为主库。</p>
     */
    @Primary
    @Bean(name = "mysqlDs")
    public DataSource mysqlDataSource() {
        return DataSourceBuilder.create()
                .url("jdbc:mysql://192.168.2.31:3306/ry?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&useSSL=true&serverTimezone=GMT%2B8")
                .username("root")
                .password("123")
                .driverClassName("com.mysql.cj.jdbc.Driver")
                .type(HikariDataSource.class)
                .build();
    }

    /**
     * 创建 OceanBase HikariCP 数据源(非 @Primary)
     *
     * <p>从 DualWriteOceanbaseProperties 中读取连接信息,
     * 包括地址、用户名、密码、连接池参数等,
     * 配置并返回一个独立的 HikariDataSource 实例。</p>
     *
     * @param props OceanBase 连接属性(自动从 Nacos 的 oceanbase.datasource 前缀注入)
     * @return 配置好的 HikariDataSource 实例,以 "oceanbaseDs" 为 Bean 名称
     */
    @Bean(name = "oceanbaseDs")
    public DataSource oceanbaseDataSource(DualWriteOceanbaseProperties props) {
        HikariConfig config = new HikariConfig();
        config.setJdbcUrl(props.getUrl());
        config.setUsername(props.getUsername());
        config.setPassword(props.getPassword());
        config.setDriverClassName(props.getDriverClassName());
        config.setMaximumPoolSize(props.getHikari().getMaximumPoolSize());
        config.setMinimumIdle(props.getHikari().getMinimumIdle());
        config.setConnectionTimeout(props.getHikari().getConnectionTimeout());
        config.setIdleTimeout(props.getHikari().getIdleTimeout());
        config.setMaxLifetime(props.getHikari().getMaxLifetime());
        config.setPoolName("OceanBasePool");
        return new HikariDataSource(config);
    }
}

注意@Primary 确保 MyBatis-Flex 和 JdbcTemplate 使用 MySQL 为主库。这是修复"MySQL 数据源未被创建"关键 Bug 的地方。

2.4 config/DualWriteAutoConfiguration.java — 条件装配入口

package com.ruoyi.system.dualwrite.config;

import com.ruoyi.system.dualwrite.executor.DualWriteExecutor;
import com.ruoyi.system.dualwrite.interceptor.DualWriteInterceptor;
import com.ruoyi.system.dualwrite.interceptor.DualWriteDdlListener;
import com.ruoyi.system.dualwrite.log.service.DualWriteLogService;
import com.ruoyi.system.dualwrite.router.DataSourceHealthIndicator;
import com.ruoyi.system.dualwrite.router.DualWriteRouter;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;

import javax.sql.DataSource;

/**
 * 双写功能自动装配入口
 *
 * <p>这是双写模块的自动化配置类,被 Spring Boot 的自动配置机制发现。
 * 当 <b>dual-write.enabled=true</b> 时,自动注册以下核心组件:</p>
 *
 * <ol>
 *   <li><b>DualWriteRouter</b> — 路由决策器,判断当前应该使用哪个数据源</li>
 *   <li><b>DualWriteExecutor</b> — 执行器,负责将 SQL 写入 OceanBase</li>
 *   <li><b>DualWriteInterceptor</b> — MyBatis 拦截器,捕获增删改操作</li>
 *   <li><b>DualWriteDdlListener</b> — DDL 监听器,同步表结构变更</li>
 *   <li><b>DataSourceHealthIndicator</b> — 健康检查器,监控数据源状态</li>
 * </ol>
 *
 * <p>当 dual-write.enabled=false(默认)时,本配置类不被加载,
 * 所有相关的 Bean 都不会创建,完全不影响原有功能。</p>
 *
 * @author ruoyi
 */
@Configuration
@EnableConfigurationProperties(DualWriteProperties.class)
@ConditionalOnProperty(prefix = "dual-write", name = "enabled", havingValue = "true", matchIfMissing = false)
@Import({AsyncConfig.class, DualWriteDataSourceConfig.class})
public class DualWriteAutoConfiguration {

    /**
     * 创建路由决策器
     * <p>根据配置的 active 角色和运行时健康状态,
     * 决定当前双写时应该以哪个数据源为主。</p>
     *
     * @param properties 双写配置
     * @param oceanbaseDs OceanBase 数据源(由 DualWriteDataSourceConfig 创建)
     * @return DualWriteRouter 实例
     */
    @Bean
    public DualWriteRouter dualWriteRouter(DualWriteProperties properties,
                                           @Qualifier("oceanbaseDs") DataSource oceanbaseDs) {
        return new DualWriteRouter(properties, oceanbaseDs);
    }

    /**
     * 创建双写执行器
     * <p>负责将拦截器捕获的 SQL 和参数在 OceanBase 上执行,
     * 支持同步和异步两种模式。</p>
     *
     * @param router 路由决策器(用于获取 OB 数据源)
     * @param logService 日志服务(用于记录双写结果)
     * @return DualWriteExecutor 实例
     */
    @Bean
    public DualWriteExecutor dualWriteExecutor(DualWriteRouter router, DualWriteLogService logService) {
        return new DualWriteExecutor(router, logService);
    }

    /**
     * 创建 MyBatis 拦截器
     * <p>拦截 MyBatis 的 Executor.update() 方法,
     * 捕获所有 INSERT/UPDATE/DELETE 操作并触发双写。</p>
     *
     * @param properties 双写配置(用于判断开关和同步模式)
     * @param executor 双写执行器
     * @param router 路由决策器
     * @return DualWriteInterceptor 实例
     */
    @Bean
    public DualWriteInterceptor dualWriteInterceptor(DualWriteProperties properties,
                                                       DualWriteExecutor executor,
                                                       DualWriteRouter router) {
        return new DualWriteInterceptor(properties, executor, router);
    }

    /**
     * 创建 DDL 同步监听器
     * <p>监听 MySQL 上的 DDL 操作,自动在 OceanBase 上执行同样的
     * CREATE/ALTER/DROP 语句,保持两库表结构一致。</p>
     *
     * @param properties 双写配置(用于判断 ddlSync 开关)
     * @param router 路由决策器(用于获取 OB 数据源)
     * @return DualWriteDdlListener 实例
     */
    @Bean
    public DualWriteDdlListener dualWriteDdlListener(DualWriteProperties properties, DualWriteRouter router) {
        return new DualWriteDdlListener(properties, router);
    }

    /**
     * 创建健康检查指示器
     * <p>通过 Actuator 的 /actuator/health 端点暴露双写数据源状态,
     * 同时内置定时任务每 10 秒检测 OB 连通性,连续 3 次失败自动切主。</p>
     *
     * @param properties 双写配置(用于判断 autoFailover 开关)
     * @param router 路由决策器
     * @param oceanbaseDs OceanBase 数据源
     * @return DataSourceHealthIndicator 实例
     */
    @Bean
    public DataSourceHealthIndicator dataSourceHealthIndicator(DualWriteProperties properties,
                                                                DualWriteRouter router,
                                                                @Qualifier("oceanbaseDs") DataSource oceanbaseDs) {
        return new DataSourceHealthIndicator(properties, router, oceanbaseDs);
    }
}

2.5 config/AsyncConfig.java — 异步线程池

package com.ruoyi.system.dualwrite.config;

import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;

import java.util.concurrent.Executor;

/**
 * 异步双写线程池配置
 *
 * <p>创建一个名为 "dualWriteExecutor" 的 Spring 线程池,
 * 用于异步执行 OceanBase 的双写操作。</p>
 *
 * <p>线程池设计考虑:</p>
 * <ul>
 *   <li>双写操作主要是 JDBC 调用,属于 IO 密集型任务,
 *       因此线程数不需要太大,默认核心 4 个、最大 8 个。</li>
 *   <li>队列容量 200,当双写请求超过线程处理能力时排队等待,
 *       避免丢弃任务。</li>
 *   <li>应用关闭时等待已提交的任务最多 30 秒,
 *       防止强制关闭导致数据丢失。</li>
 *   <li>线程名前缀 "dual-write-" 方便通过 jstack 等工具定位排查。</li>
 * </ul>
 *
 * @author ruoyi
 */
@Configuration
@EnableAsync
@ConditionalOnProperty(prefix = "dual-write", name = "enabled", havingValue = "true", matchIfMissing = false)
public class AsyncConfig {

    /**
     * 创建并配置双写专用线程池
     *
     * <p>参数通过 DualWriteProperties.threadPool 注入,
     * 如果不在配置文件中显式指定则使用默认值。</p>
     *
     * @param props 双写配置(从中读取 threadPool.corePoolSize 等)
     * @return 配置好的 ThreadPoolTaskExecutor 实例
     */
    @Bean(name = "dualWriteTaskExecutor")
    public Executor dualWriteTaskExecutor(DualWriteProperties props) {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(props.getThreadPool().getCorePoolSize());
        executor.setMaxPoolSize(props.getThreadPool().getMaxPoolSize());
        executor.setQueueCapacity(props.getThreadPool().getQueueCapacity());
        executor.setThreadNamePrefix("dual-write-");
        // 优雅关闭:等待已提交但尚未完成的任务执行完毕
        executor.setWaitForTasksToCompleteOnShutdown(true);
        executor.setAwaitTerminationSeconds(30);
        executor.initialize();
        return executor;
    }
}

2.6 executor/DualWriteTask.java — 任务载体

package com.ruoyi.system.dualwrite.executor;

import org.apache.ibatis.mapping.MappedStatement;

/**
 * 双写任务封装
 *
 * <p>将 MyBatis 拦截器捕获的 INSERT / UPDATE / DELETE 信息
 * 封装为一个可以在线程池中独立执行的任务单元。</p>
 *
 * <p>包含的信息:</p>
 * <ul>
 *   <li>MappedStatement — MyBatis 的 SQL 元信息(SQL 语句、参数映射等)</li>
 *   <li>parameter — 实体对象或 Map 类型的参数</li>
 *   <li>tableName — 操作的表名(用于日志和监控)</li>
 *   <li>operation — 操作类型(INSERT / UPDATE / DELETE)</li>
 * </ul>
 *
 * <p>DualWriteExecutor 接收到 DualWriteTask 后,从中提取 SQL 和参数,
 * 然后在 OceanBase 上执行相同的操作。</p>
 *
 * @author ruoyi
 */
public class DualWriteTask {

    /**
     * MyBatis MappedStatement
     * <p>包含了完整的 SQL 语句、参数类型、返回值类型等元信息。
     * 通过 getBoundSql(parameter) 可以获取最终的 SQL 和参数映射列表。</p>
     */
    private final MappedStatement mappedStatement;

    /**
     * SQL 执行参数
     * <p>可以是实体对象(如 SysUser)、Map 或简单类型。
     * 由 MappedStatement 中的参数映射定义决定如何从该对象中提取值。</p>
     */
    private final Object parameter;

    /**
     * 被操作的表名
     * <p>由 DualWriteInterceptor 从 SQL 语句中解析提取,
     * 用于记录日志和监控统计。</p>
     */
    private final String tableName;

    /**
     * 操作类型
     * <p>INSERT / UPDATE / DELETE 三种之一,
     * 用于日志分类和后续补偿处理。</p>
     */
    private final String operation;

    /**
     * 构造双写任务
     *
     * @param mappedStatement MyBatis SQL 元信息
     * @param parameter       SQL 参数
     * @param tableName       操作表名
     * @param operation       操作类型
     */
    public DualWriteTask(MappedStatement mappedStatement, Object parameter,
                          String tableName, String operation) {
        this.mappedStatement = mappedStatement;
        this.parameter = parameter;
        this.tableName = tableName;
        this.operation = operation;
    }

    public MappedStatement getMappedStatement() { return mappedStatement; }
    public Object getParameter() { return parameter; }
    public String getTableName() { return tableName; }
    public String getOperation() { return operation; }
}

2.7 executor/DualWriteExecutor.java — 双写执行器(核心)

package com.ruoyi.system.dualwrite.executor;

import com.ruoyi.system.dualwrite.log.service.DualWriteLogService;
import com.ruoyi.system.dualwrite.router.DualWriteRouter;
import org.apache.ibatis.mapping.BoundSql;
import org.apache.ibatis.mapping.MappedStatement;
import org.apache.ibatis.mapping.ParameterMapping;
import org.apache.ibatis.reflection.MetaObject;
import org.apache.ibatis.session.Configuration;
import org.apache.ibatis.type.TypeHandlerRegistry;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.scheduling.annotation.Async;

import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.text.DateFormat;
import java.util.Date;
import java.util.List;
import java.util.Locale;

/**
 * 双写执行器
 *
 * <p>核心职责:将 DualWriteTask 中的 SQL + 参数在 OceanBase 上执行。</p>
 *
 * <p>执行模式:</p>
 * <ul>
 *   <li><b>异步(executeAsync)</b> — 通过 @Async 提交到 dualWriteExecutor 线程池,
 *       不阻塞主库请求。适合对性能要求高、可接受短暂数据不一致的场景。</li>
 *   <li><b>同步(executeSync)</b> — 在当前线程中执行,失败时抛出异常,
 *       会让 MyBatis 拦截器收到异常,最终导致整个事务回滚。
 *       适合对数据一致性要求极高的场景。</li>
 * </ul>
 *
 * <p>参数处理:</p>
 * <ul>
 *   <li>通过 MyBatis 的 BoundSql 获取带 ? 占位符的 SQL</li>
 *   <li>从 ParameterMapping 中逐个获取参数值</li>
 *   <li>Date 类型按中文格式转为字符串,避免时区问题</li>
 *   <li>其他类型通过 JDBC 的 setObject 自动匹配</li>
 * </ul>
 *
 * @author ruoyi
 */
public class DualWriteExecutor {

    private static final Logger log = LoggerFactory.getLogger(DualWriteExecutor.class);

    private final DualWriteRouter router;
    private final DualWriteLogService logService;

    public DualWriteExecutor(DualWriteRouter router, DualWriteLogService logService) {
        this.router = router;
        this.logService = logService;
    }

    /**
     * 异步执行双写
     *
     * <p>被 @Async("dualWriteExecutor") 注解标记,Spring 会自动将调用
     * 提交到 dualWriteExecutor 线程池中异步执行。
     * 调用方(拦截器)不会等待结果。</p>
     *
     * <p>异常处理:任何异常都会被捕获并记录到 dual_write_log 表,
     * 不会抛出给调用方。</p>
     *
     * @param task 双写任务(包含 SQL 和参数)
     */
@Async("dualWriteTaskExecutor")
    public void executeAsync(DualWriteTask task) {
        try {
            execute(task);
        } catch (Exception e) {
            log.error("异步双写失败,表:{},操作:{},错误:{}",
                    task.getTableName(), task.getOperation(), e.getMessage());
            saveLog(task, "FAILED", e.getMessage());
        }
    }

    /**
     * 同步执行双写
     *
     * <p>在当前调用线程中直接执行,不经过线程池。
     * 如果 OceanBase 写入失败,异常会抛出给调用方(DualWriteInterceptor),
     * 导致整个数据库事务回滚。</p>
     *
     * @param task 双写任务(包含 SQL 和参数)
     * @throws Exception 执行失败时抛出,由拦截器决定是否回滚事务
     */
    public void executeSync(DualWriteTask task) throws Exception {
        try {
            execute(task);
        } catch (Exception e) {
            log.error("同步双写失败,表:{},操作:{},错误:{}",
                    task.getTableName(), task.getOperation(), e.getMessage());
            saveLog(task, "FAILED", e.getMessage());
            throw e;
        }
    }

    /**
     * 核心执行逻辑
     *
     * <p>步骤:</p>
     * <ol>
     *   <li>从 MappedStatement 获取 BoundSql(含 SQL 语句和参数映射)</li>
     *   <li>从 BoundSql 获取最终 SQL 语句(带 ? 占位符)</li>
     *   <li>获取 OceanBase 数据源连接</li>
     *   <li>创建 PreparedStatement,设置参数</li>
     *   <li>执行 executeUpdate(),写入 OceanBase</li>
     *   <li>记录成功日志</li>
     * </ol>
     *
     * @param task 双写任务
     * @throws Exception SQL 执行异常或连接异常
     */
    private void execute(DualWriteTask task) throws Exception {
        MappedStatement ms = task.getMappedStatement();
        Object parameter = task.getParameter();
        BoundSql boundSql = ms.getBoundSql(parameter);
        String sql = boundSql.getSql();

        DataSource obDs = router.getOceanbaseDs();
        if (obDs == null) {
            log.warn("OceanBase 数据源不可用,跳过双写(表:{})", task.getTableName());
            return;
        }

        try (Connection conn = obDs.getConnection();
             PreparedStatement ps = conn.prepareStatement(sql)) {
            setParameters(ps, ms, boundSql, parameter);
            int rows = ps.executeUpdate();
            saveLog(task, "SUCCESS", null);
            log.debug("双写成功,表:{},影响行数:{}", task.getTableName(), rows);
        }
    }

    /**
     * 将 MyBatis 参数映射设置到 JDBC PreparedStatement
     *
     * <p>MyBatis 的 BoundSql 中保存了参数映射列表(ParameterMapping),
     * 每个映射包含参数名、类型等信息。本方法遍历该列表,
     * 从参数对象中提取对应的值,设置到 PreparedStatement 的位置参数上。</p>
     *
     * <p>特殊处理:</p>
     * <ul>
     *   <li>Date 类型 → 转为 "yyyy-MM-dd HH:mm:ss" 格式的中文字符串
     *       (避免 OceanBase 的时区和 MyBatis 的 TypeHandler 不一致)</li>
     *   <li>__frch_ 前缀 → MyBatis 的 foreach 动态 SQL 生成的前缀</li>
     *   <li>additionalParameter → MyBatis 中通过 setAdditionalParameter 设置的额外参数</li>
     * </ul>
     *
     * @param ps        目标 PreparedStatement
     * @param ms        MyBatis MappedStatement
     * @param boundSql  绑定的 SQL(含参数映射)
     * @param parameter 参数对象
     * @throws Exception 参数设置异常
     */
    private void setParameters(PreparedStatement ps, MappedStatement ms,
                                BoundSql boundSql, Object parameter) throws Exception {
        Configuration configuration = ms.getConfiguration();
        List<ParameterMapping> parameterMappings = boundSql.getParameterMappings();
        if (parameterMappings == null || parameterMappings.isEmpty()) {
            return;
        }

        MetaObject metaObject = configuration.newMetaObject(parameter);
        TypeHandlerRegistry typeHandlerRegistry = configuration.getTypeHandlerRegistry();

        for (int i = 0; i < parameterMappings.size(); i++) {
            ParameterMapping pm = parameterMappings.get(i);
            String propertyName = pm.getProperty();

            Object value;
            // 判断参数类型:如果是简单类型直接使用,否则从对象中获取属性值
            if (typeHandlerRegistry.hasTypeHandler(parameter.getClass())) {
                value = parameter;
            } else if (boundSql.hasAdditionalParameter(propertyName)) {
                value = boundSql.getAdditionalParameter(propertyName);
            } else if (propertyName.startsWith("__frch_")) {
                // foreach 动态 SQL 生成的参数前缀
                value = metaObject.getValue(propertyName);
            } else {
                value = metaObject.getValue(propertyName);
            }

            // Date 类型特殊处理:避免 OB 时区问题
            if (value instanceof Date) {
                ps.setString(i + 1,
                    DateFormat.getDateTimeInstance(
                        DateFormat.DEFAULT, DateFormat.DEFAULT, Locale.CHINA
                    ).format(value));
            } else {
                ps.setObject(i + 1, value);
            }
        }
    }

    /**
     * 记录双写日志到 dual_write_log 表
     *
     * @param task      双写任务信息
     * @param status    执行状态(SUCCESS / FAILED)
     * @param errorMsg  错误信息(成功时为 null)
     */
    private void saveLog(DualWriteTask task, String status, String errorMsg) {
        try {
            logService.saveLog(task.getTableName(), task.getOperation(),
                    "OceanBase", status, errorMsg);
        } catch (Exception e) {
            log.warn("记录双写日志失败,表:{},错误:{}", task.getTableName(), e.getMessage());
        }
    }
}

2.8 executor/SnowflakeIdGenerator.java — 雪花算法

package com.ruoyi.system.dualwrite.executor;

/**
 * 雪花算法 ID 生成器
 *
 * <p>为什么需要雪花算法?</p>
 * <ul>
 *   <li>MySQL 使用自增 ID(auto_increment),OceanBase 也使用自增 ID</li>
 *   <li>双写时,MySQL 生成的 ID 和 OceanBase 自增的 ID 不一致</li>
 *   <li>如果直接用自增 ID,两个库的同一条记录会有不同的主键值</li>
 *   <li>使用雪花算法在应用层生成全局唯一 ID,代替数据库自增</li>
 *   <li>这样 MySQL 和 OceanBase 中同一条记录的主键完全相同</li>
 * </ul>
 *
 * <p>ID 结构(64 位 long):</p>
 * <pre>
 *  0 | 0000000000 0000000000 0000000000 0000000000 0 | 00000 | 00000 | 000000000000
 *  ↑                        ↑                          ↑       ↑           ↑
 * 符号位                 时间戳(41位)              数据中心  工作节点    序列号
 * (始终为0)           (毫秒级,差值)             ID(5位)  ID(5位)   (12位,4096/毫秒)
 * </pre>
 *
 * <p>性能:单机理论 QPS 409.6 万/秒,ID 趋势递增。</p>
 *
 * @author ruoyi
 */
public class SnowflakeIdGenerator {

    /** 工作节点 ID(0~31),每个实例需要不同 */
    private final long workerId;

    /** 数据中心 ID(0~31) */
    private final long datacenterId;

    /**
     * 起始时间戳
     * <p>取 2023-11-15 的毫秒值,可以根据实际上线日期调整。
     * 该值越小,ID 可用时间越长(可用 69 年)。</p>
     */
    private final long epoch = 1700000000000L;

    /** 工作节点 ID 位数(5位,取值范围 0~31) */
    private final long workerIdBits = 5L;

    /** 数据中心 ID 位数(5位,取值范围 0~31) */
    private final long datacenterIdBits = 5L;

    /** 序列号位数(12位,每毫秒最多 4096 个 ID) */
    private final long sequenceBits = 12L;

    /** 工作节点 ID 左移位数 = 序列号位数(12) */
    private final long workerIdShift = sequenceBits;

    /** 数据中心 ID 左移位数 = 12 + 5 = 17 */
    private final long datacenterIdShift = sequenceBits + workerIdBits;

    /** 时间戳左移位数 = 12 + 5 + 5 = 22 */
    private final long timestampLeftShift = sequenceBits + workerIdBits + datacenterIdBits;

    /** 序列号掩码 = 2^12 - 1 = 4095 */
    private final long sequenceMask = -1L ^ (-1L << sequenceBits);

    /** 上次生成 ID 的时间戳(用于检测时钟回拨和毫秒切换) */
    private long lastTimestamp = -1L;

    /** 当前毫秒内的序列号(0~4095,溢出后等待下一毫秒) */
    private long sequence = 0L;

    /**
     * 构造雪花 ID 生成器
     *
     * @param workerId     工作节点 ID(0~31),同一数据中心内不能重复
     * @param datacenterId 数据中心 ID(0~31),跨数据中心时不能重复
     * @throws IllegalArgumentException 参数超出范围时抛出
     */
    public SnowflakeIdGenerator(long workerId, long datacenterId) {
        if (workerId > 31 || workerId < 0) {
            throw new IllegalArgumentException("workerId 必须在 0~31 之间");
        }
        if (datacenterId > 31 || datacenterId < 0) {
            throw new IllegalArgumentException("datacenterId 必须在 0~31 之间");
        }
        this.workerId = workerId;
        this.datacenterId = datacenterId;
    }

    /**
     * 生成下一个全局唯一 ID(线程安全)
     *
     * <p>同步方法,同一时刻只有一个线程可以生成 ID。
     * 每毫秒最多支持 4096 个唯一 ID,超出后自动等待下一毫秒。</p>
     *
     * <p>时钟回拨处理:</p>
     * <ul>
     *   <li>检测到系统时间回拨时,不抛出异常,而是使用上次的时间戳</li>
     *   <li>这样会导致序列号继续累加,但不会产生 ID 冲突</li>
     *   <li>如果是大幅度回拨(跨天级别),建议重启应用</li>
     * </ul>
     *
     * @return 64 位长整型全局唯一 ID
     */
    public synchronized long nextId() {
        long timestamp = System.currentTimeMillis();

        // 时钟回拨保护:如果当前时间小于上次生成时间,用上次时间戳兜底
        if (timestamp < lastTimestamp) {
            timestamp = lastTimestamp;
        }

        // 同一毫秒内,序列号递增并检查溢出
        if (timestamp == lastTimestamp) {
            sequence = (sequence + 1) & sequenceMask;
            if (sequence == 0) {
                // 序列号用尽(4095 已用完),等待下一毫秒
                timestamp = tilNextMillis(lastTimestamp);
            }
        } else {
            // 进入新的一毫秒,序列号重置为 0
            sequence = 0L;
        }

        lastTimestamp = timestamp;

        // 拼接各部分生成最终 ID
        return ((timestamp - epoch) << timestampLeftShift)
                | (datacenterId << datacenterIdShift)
                | (workerId << workerIdShift)
                | sequence;
    }

    /**
     * 自旋等待到下一毫秒
     *
     * <p>当同一毫秒内序列号用完时,循环等待直到系统时间进入下一毫秒。</p>
     *
     * @param lastTimestamp 上次生成 ID 的时间戳(毫秒)
     * @return 下一毫秒的时间戳
     */
    private long tilNextMillis(long lastTimestamp) {
        long timestamp = System.currentTimeMillis();
        while (timestamp <= lastTimestamp) {
            timestamp = System.currentTimeMillis();
        }
        return timestamp;
    }
}

2.9 interceptor/DualWriteInterceptor.java — 核心拦截器

package com.ruoyi.system.dualwrite.interceptor;

import com.ruoyi.system.dualwrite.config.DualWriteProperties;
import com.ruoyi.system.dualwrite.executor.DualWriteExecutor;
import com.ruoyi.system.dualwrite.executor.DualWriteTask;
import com.ruoyi.system.dualwrite.router.DualWriteRouter;
import org.apache.ibatis.executor.Executor;
import org.apache.ibatis.mapping.MappedStatement;
import org.apache.ibatis.mapping.SqlCommandType;
import org.apache.ibatis.plugin.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.sql.Connection;
import java.util.Properties;

/**
 * MyBatis Executor 拦截器 —— 双写功能的核心入口
 *
 * <p>拦截 MyBatis 的 {@link Executor#update(MappedStatement, Object)} 方法,
 * 在 MyBatis 执行完 INSERT / UPDATE / DELETE 后,
 * 将同样的 SQL + 参数同步写入 OceanBase。</p>
 *
 * <p>工作流程:</p>
 * <ol>
 *   <li>检查 dual-write.enabled 开关,关闭时直接放行</li>
 *   <li>调用 invocation.proceed() 先让 MyBatis 正常执行主库(MySQL)的 SQL</li>
 *   <li>获取 SqlCommandType,只处理 INSERT / UPDATE / DELETE</li>
 *   <li>通过关键词从 SQL 中提取表名,排除 dual_write_log 表(避免递归死循环)</li>
 *   <li>根据 syncMode 选择同步或异步方式将 SQL 写入 OceanBase</li>
 * </ol>
 *
 * <p>排除 dual_write_log 的原因:</p>
 * <ul>
 *   <li>DualWriteExecutor 中会调用 logService.saveLog() 写入 dual_write_log 表</li>
 *   <li>如果不排除该表,写入日志时会再次触发拦截器,形成无限递归</li>
 *   <li>日志表只存在于 MySQL,不需要双写到 OB</li>
 * </ul>
 *
 * @author ruoyi
 */
@Intercepts({
    @Signature(type = Executor.class, method = "update", args = {MappedStatement.class, Object.class})
})
public class DualWriteInterceptor implements Interceptor {

    private static final Logger log = LoggerFactory.getLogger(DualWriteInterceptor.class);

    private final DualWriteProperties properties;
    private final DualWriteExecutor executor;
    @SuppressWarnings("unused")
    private final DualWriteRouter router;

    public DualWriteInterceptor(DualWriteProperties properties,
                                 DualWriteExecutor executor,
                                 DualWriteRouter router) {
        this.properties = properties;
        this.executor = executor;
        this.router = router;
    }

    /**
     * 拦截方法:MyBatis 每次执行 update(包括 insert/update/delete)时调用
     *
     * <p>执行顺序:</p>
     * <ol>
     *   <li>首先调用 invocation.proceed() 执行原始的 SQL(写入 MySQL)</li>
     *   <li>获取执行结果(影响行数)</li>
     *   <li>根据 SQL 类型和表名决定是否需要双写</li>
     *   <li>按照配置的同步模式执行双写</li>
     *   <li>返回主库的执行结果</li>
     * </ol>
     *
     * @param invocation MyBatis 调用信息,包含 MappedStatement 和参数
     * @return 主库 SQL 执行结果(影响行数)
     * @throws Throwable 主库执行异常时直接抛出,双写异常视模式而定
     */
    @Override
    public Object intercept(Invocation invocation) throws Throwable {
        // 双写总开关:为 false 时完全不影响 MyBatis 行为
        if (!properties.isEnabled()) {
            return invocation.proceed();
        }

        Object[] args = invocation.getArgs();
        MappedStatement ms = (MappedStatement) args[0];
        Object parameter = args[1];
        SqlCommandType sqlType = ms.getSqlCommandType();

        // 第一步:先执行主库(MySQL),获取执行结果
        log.warn("=== DUALWRITE === SQL: {}", ms.getBoundSql(parameter).getSql().replace("\r"," ").replace("\n"," ").replaceAll("  +"," "));
        log.warn("=== DUALWRITE === type: {}, table: {}, param: {}", sqlType, extractTableName(ms), parameter);
        Object result = invocation.proceed();
        log.warn("=== DUALWRITE === result: {} rows", result);

        // 第二步:判断 SQL 类型,只有写操作需要双写
        if (sqlType == SqlCommandType.INSERT || sqlType == SqlCommandType.UPDATE
                || sqlType == SqlCommandType.DELETE) {
            String tableName = extractTableName(ms);
            // 排除日志表自身,否则 saveLog -> 插入 dual_write_log -> 再次触发拦截器 -> 死循环
            if (tableName != null && !"dual_write_log".equals(tableName)) {
                DualWriteTask task = new DualWriteTask(ms, parameter, tableName, sqlType.name());
                String syncMode = properties.getSyncMode();
                if ("SYNC".equalsIgnoreCase(syncMode)) {
                    // 同步模式:等待 OB 写入完成,写入失败时抛异常,整个事务回滚
                    log.debug("同步双写表:{}", tableName);
                    executor.executeSync(task);
                } else {
                    // 异步模式:提交到线程池后立即返回,不阻塞主库响应
                    log.debug("异步双写表:{}", tableName);
                    try {
                        executor.executeAsync(task);
                    } catch (Exception e) {
                        // 异步提交异常时不能影响主库事务
                        log.error("异步双写提交异常,表:{},错误:{}", tableName, e.getMessage());
                    }
                }
            }
        }

        return result;
    }

    /**
     * 从 MappedStatement 的 SQL 中解析表名
     *
     * <p>通过匹配 INSERT INTO / UPDATE / DELETE FROM 等关键词来定位表名位置。
     * 这种方法不依赖 XML 映射配置,即使是用 MyBatis-Flex 自动生成的 SQL 也能正确解析。</p>
     *
     * <p>解析示例:</p>
     * <ul>
     *   <li>"INSERT INTO sys_user (id, name) VALUES (?, ?)" → "sys_user"</li>
     *   <li>"UPDATE sys_config SET value=? WHERE id=?" → "sys_config"</li>
     *   <li>"DELETE FROM sys_dept WHERE id=?" → "sys_dept"</li>
     * </ul>
     *
     * @param ms MyBatis MappedStatement
     * @return 解析出的表名(已去除反引号和引号),解析失败返回 null
     */
    private String extractTableName(MappedStatement ms) {
        String sql = ms.getBoundSql(null).getSql();
        sql = sql.replace("\n", " ").replace("\r", " ").trim();
        String upper = sql.toUpperCase();
        String[] keywords = {"INSERT INTO", "UPDATE ", "DELETE FROM"};
        for (String kw : keywords) {
            int idx = upper.indexOf(kw);
            if (idx >= 0) {
                String after = sql.substring(idx + kw.length()).trim();
                int end = after.indexOf(' ');
                if (end > 0) {
                    return after.substring(0, end).replace("`", "").replace("\"", "");
                }
                return after.replace("`", "").replace("\"", "");
            }
        }
        return null;
    }

    /**
     * 包装目标对象
     *
     * @param target 被拦截的目标对象(Executor 实例)
     * @return 代理对象
     */
    @Override
    public Object plugin(Object target) {
        return Plugin.wrap(target, this);
    }

    /**
     * 设置 MyBatis 插件属性(本拦截器不需要额外属性)
     */
    @Override
    public void setProperties(Properties properties) {
    }
}

2.10 interceptor/DualWriteDdlListener.java — DDL 同步

package com.ruoyi.system.dualwrite.interceptor;

import com.ruoyi.system.dualwrite.config.DualWriteProperties;
import com.ruoyi.system.dualwrite.router.DualWriteRouter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.Statement;

/**
 * DDL 同步监听器
 *
 * <p>负责将 MySQL 上的 DDL 语句(CREATE TABLE / ALTER TABLE / DROP TABLE 等)
 * 同步执行到 OceanBase,确保两库的表结构始终一致。</p>
 *
 * <p>触发方式:</p>
 * <ul>
 *   <li>由业务代码在 MyBatis Mapper 的 XML 中执行 DDL 后手动调用 onDdl()</li>
 *   <li>例如在 GenTableMapper 中执行 CREATE TABLE 后调用本监听器</li>
 * </ul>
 *
 * <p>注意:</p>
 * <ul>
 *   <li>DDL 是数据库层面的操作,MyBatis 的 Executor 拦截器只能拦截 DML,不能拦截 DDL</li>
 *   <li>因此 DDL 同步需要业务代码主动调用,无法自动拦截</li>
 *   <li>DDL 同步默认开启,可通过 dual-write.ddl-sync=false 关闭</li>
 * </ul>
 *
 * @author ruoyi
 */
public class DualWriteDdlListener {

    private static final Logger log = LoggerFactory.getLogger(DualWriteDdlListener.class);

    private final DualWriteProperties properties;
    private final DualWriteRouter router;

    public DualWriteDdlListener(DualWriteProperties properties, DualWriteRouter router) {
        this.properties = properties;
        this.router = router;
    }

    /**
     * 执行 DDL 同步
     *
     * <p>将指定的 DDL 语句在 OceanBase 上执行一次。
     * 如果 OB 上已经存在同名表,CREATE TABLE 会报错(与 MySQL 行为一致)。</p>
     *
     * <p>使用场景示例:</p>
     * <pre>
     * // 代码生成器创建新表后,同步到 OceanBase
     * ddlListener.onDdl("CREATE TABLE `gen_table_xxx` (...)");
     *
     * // 升级脚本执行 ALTER TABLE 后同步
     * ddlListener.onDdl("ALTER TABLE sys_user ADD COLUMN ...");
     * </pre>
     *
     * @param ddlSql 完整的 DDL 语句,如 "CREATE TABLE sys_xxx (...)"
     */
    public void onDdl(String ddlSql) {
        if (!properties.isEnabled() || !properties.isDdlSync()) {
            return;
        }
        DataSource obDs = router.getOceanbaseDs();
        if (obDs == null) {
            log.warn("OceanBase 数据源不可用,跳过 DDL 同步");
            return;
        }
        try (Connection conn = obDs.getConnection();
             Statement stmt = conn.createStatement()) {
            stmt.execute(ddlSql);
            log.info("DDL 已同步到 OceanBase:{}", ddlSql);
        } catch (Exception e) {
            log.error("DDL 同步到 OceanBase 失败:{},错误:{}", ddlSql, e.getMessage());
        }
    }
}

2.11 router/DualWriteRouter.java — 路由决策器

package com.ruoyi.system.dualwrite.router;

import com.ruoyi.system.dualwrite.config.DualWriteProperties;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.Statement;

/**
 * 双写路由决策器
 *
 * <p>核心职责:</p>
 * <ol>
 *   <li><b>角色管理</b> — 根据配置和运行时健康状态,维护当前主库角色</li>
 *   <li><b>故障切换</b> — 提供 failover/recover 方法,支持主从互换</li>
 *   <li><b>连通性检测</b> — 提供 isOceanbaseAvailable() 方法检测 OB 状态</li>
 * </ol>
 *
 * <p>线程安全性说明:</p>
 * <ul>
 *   <li>primaryRole 和 activeRole 使用 volatile 关键字声明,
 *       确保一个线程修改后其他线程立即可见。</li>
 *   <li>failover/recover 是轻量操作,仅修改两个枚举变量,
 *       不会产生并发问题。</li>
 * </ul>
 *
 * @author ruoyi
 */
public class DualWriteRouter {

    private static final Logger log = LoggerFactory.getLogger(DualWriteRouter.class);

    /** 双写配置(用于读取 active 初始值) */
    private final DualWriteProperties properties;

    /** OceanBase 数据源引用(由外部注入,始终指向同一个 OB 连接池) */
    private final DataSource oceanbaseDs;

    /**
     * 主库角色
     * <p>volatile 保证 failover 切换后所有线程立即可见新的角色值。
     * MASTER = MySQL 为主,SLAVE = OceanBase 为主。</p>
     */
    private volatile DataSourceRole primaryRole = DataSourceRole.MASTER;

    /**
     * 当前活跃角色
     * <p>在 failover 过程中,primaryRole 改变后 activeRole 也随之改变。
     * 将两者分开是为了将来可能的灰度切换场景。</p>
     */
    private volatile DataSourceRole activeRole = DataSourceRole.MASTER;

    /**
     * 构造路由决策器
     *
     * <p>如果配置中 active=SLAVE,则初始化时将主库设为 OceanBase,
     * 实现从"MySQL+OB 双写、OB 为主"的模式启动。</p>
     *
     * @param properties  双写配置(读取 active 初始值)
     * @param oceanbaseDs OceanBase 数据源
     */
    public DualWriteRouter(DualWriteProperties properties, DataSource oceanbaseDs) {
        this.properties = properties;
        this.oceanbaseDs = oceanbaseDs;
        String active = properties.getActive();
        if ("SLAVE".equalsIgnoreCase(active)) {
            this.activeRole = DataSourceRole.SLAVE;
            this.primaryRole = DataSourceRole.SLAVE;
        }
    }

    /**
     * 获取当前主库数据源
     * <p>如果主库是 MySQL,返回 null(MySQL 数据源不由本模块管理);
     * 如果主库已切换到 OceanBase,返回 OB 数据源。</p>
     *
     * @return OceanBase 数据源或 null
     */
    public DataSource getPrimaryDs() {
        return primaryRole == DataSourceRole.MASTER ? null : oceanbaseDs;
    }

    /**
     * 获取当前从库数据源
     * <p>与 getPrimaryDs() 相反,主库是 MySQL 时返回 OB,反之返回 null。</p>
     *
     * @return OceanBase 数据源或 null
     */
    public DataSource getSecondaryDs() {
        return primaryRole == DataSourceRole.MASTER ? oceanbaseDs : null;
    }

    /**
     * 获取当前活跃数据源
     * <p>与 primaryRole 的区别在于,activeRole 可能用于更细粒度的控制,
     * 当前实现与 primaryRole 保持同步变化。</p>
     *
     * @return 当前活跃数据源或 null(为 MySQL 时返回 null)
     */
    public DataSource getActiveDs() {
        return activeRole == DataSourceRole.MASTER ? null : oceanbaseDs;
    }

    /**
     * 获取 OceanBase 数据源
     * <p>此方法始终返回 OB 数据源引用,无论当前角色是什么。
     * 用于双写执行器直接获取 OB 连接进行写入。</p>
     *
     * @return OceanBase 数据源(始终非空)
     */
    public DataSource getOceanbaseDs() {
        return oceanbaseDs;
    }

    /**
     * 判断当前是否以 MySQL 为主库
     *
     * @return true = MySQL 是主库,false = OceanBase 已升级为主库
     */
    public boolean isMasterActive() {
        return activeRole == DataSourceRole.MASTER;
    }

    /**
     * 故障切换:将主库从 MySQL 切换到 OceanBase
     *
     * <p>当 MySQL 连续健康检查失败超过阈值时调用此方法。
     * 切换后新的写入会发往 OceanBase,MySQL 恢复后不再承担主库职责。</p>
     */
    public void failover() {
        if (primaryRole == DataSourceRole.MASTER) {
            primaryRole = DataSourceRole.SLAVE;
            activeRole = DataSourceRole.SLAVE;
            log.warn("【故障切换】主库已从 MySQL 切换至 OceanBase");
        }
    }

    /**
     * 恢复:将主库从 OceanBase 切换回 MySQL
     *
     * <p>当 OceanBase 健康检查恢复且连续成功时调用此方法。
     * 切换后主库职责交还给 MySQL,OceanBase 回到从库角色。</p>
     */
    public void recover() {
        if (primaryRole == DataSourceRole.SLAVE) {
            primaryRole = DataSourceRole.MASTER;
            activeRole = DataSourceRole.MASTER;
            log.info("【恢复】主库已从 OceanBase 恢复至 MySQL");
        }
    }

    /**
     * 检测 OceanBase 是否可用
     *
     * <p>通过建立 JDBC 连接并执行 SELECT 1 来验证 OB 是否可达。
     * 每次检测都是实时连接,不依赖缓存的状态。</p>
     *
     * @return true = OB 连接正常可读写,false = OB 不可用
     */
    public boolean isOceanbaseAvailable() {
        try (Connection conn = oceanbaseDs.getConnection();
             Statement stmt = conn.createStatement()) {
            stmt.execute("SELECT 1");
            return true;
        } catch (Exception e) {
            return false;
        }
    }
}

2.12 router/DataSourceRole.java — 角色枚举

package com.ruoyi.system.dualwrite.router;

/**
 * 数据源角色枚举
 *
 * <p>定义双写架构中两个数据源的角色:</p>
 * <ul>
 *   <li><b>MASTER</b> — 主库,默认是 MySQL,负责处理所有的读写请求。
 *       在主库正常时,所有业务操作都走 MySQL。</li>
 *   <li><b>SLAVE</b> — 从库,默认是 OceanBase,作为 MySQL 的数据备份。
 *       正常情况下只接收双写的写入,不承担读请求。
 *       当 MySQL 发生故障且 failover 触发后,SLAVE 会升级为新的主库。</li>
 * </ul>
 *
 * <p>角色切换流程:</p>
 * <ol>
 *   <li>正常状态:MASTER = MySQL,SLAVE = OceanBase</li>
 *   <li>MySQL 故障:触发 failover(),SLAVE 升级为 MASTER</li>
 *   <li>MySQL 恢复:触发 recover(),MASTER 切回 MySQL</li>
 * </ol>
 *
 * @author ruoyi
 */
public enum DataSourceRole {
    /** 主库角色(默认 MySQL) */
    MASTER,
    /** 从库角色(默认 OceanBase) */
    SLAVE
}

2.13 router/DataSourceHealthIndicator.java — 健康检查 + 自动切换

package com.ruoyi.system.dualwrite.router;

import com.ruoyi.system.dualwrite.config.DualWriteProperties;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.health.contributor.Health;
import org.springframework.boot.health.contributor.HealthIndicator;
import org.springframework.scheduling.annotation.Scheduled;

import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.Statement;

/**
 * OceanBase 健康检查指示器
 *
 * <p>双实现:</p>
 * <ul>
 *   <li><b>HealthIndicator 接口</b> — 通过 Spring Boot Actuator 的
 *       /actuator/health 端点暴露双写数据源的健康状态,
 *       方便监控系统(如 Prometheus + Grafana)采集。</li>
 *   <li><b>@Scheduled 定时任务</b> — 每 10 秒自动检测一次 OceanBase 连通性,
 *       连续 3 次失败后触发 failover。</li>
 * </ul>
 *
 * <p>故障切换流程:</p>
 * <ol>
 *   <li>健康检查检测到 OB 连接失败,consecutiveFailures +1</li>
 *   <li>连续失败 3 次后,log 输出告警,调用 router.failover()</li>
 *   <li>主库切换到 OceanBase,后续请求走 OB</li>
 *   <li>OB 恢复后连续检测成功,consecutiveFailures 重置为 0</li>
 *   <li>调用 router.recover() 将主库切回 MySQL</li>
 * </ol>
 *
 * @author ruoyi
 */
public class DataSourceHealthIndicator implements HealthIndicator {

    private static final Logger log = LoggerFactory.getLogger(DataSourceHealthIndicator.class);

    private final DualWriteProperties properties;
    private final DualWriteRouter router;
    private final DataSource oceanbaseDs;

    /** 当前连续失败次数,每 10 秒检测一次时累加 */
    private int consecutiveFailures = 0;

    /** 触发自动切换的阈值,连续失败达到此次数后执行 failover */
    private final int maxFailures = 3;

    /** OceanBase 是否健康,供其他组件快速判断而不必建立 JDBC 连接 */
    private volatile boolean oceanbaseHealthy = true;

    public DataSourceHealthIndicator(DualWriteProperties properties, DualWriteRouter router, DataSource oceanbaseDs) {
        this.properties = properties;
        this.router = router;
        this.oceanbaseDs = oceanbaseDs;
    }

    /**
     * Actuator 健康检查端点
     *
     * <p>被 Spring Boot Actuator 框架自动调用,返回当前双写数据源的健康状态。
     * 访问 /actuator/health 可以看到类似以下内容:</p>
     * <pre>
     * {
     *   "status": "UP",
     *   "components": {
     *     "oceanbase": { "status": "UP" },
     *     "currentPrimary": "MySQL"
     *   }
     * }
     * </pre>
     *
     * @return Health 对象,包含 OB 状态和当前主库信息
     */
    @Override
    public Health health() {
        boolean obOk = checkOceanbase();
        oceanbaseHealthy = obOk;
        if (obOk) {
            return Health.up()
                    .withDetail("oceanbase", "UP")
                    .withDetail("currentPrimary", router.isMasterActive() ? "MySQL" : "OceanBase")
                    .build();
        }
        return Health.down()
                .withDetail("oceanbase", "DOWN")
                .withDetail("currentPrimary", router.isMasterActive() ? "MySQL" : "OceanBase")
                .build();
    }

    /**
     * 定时健康检查(每 10 秒执行一次)
     *
     * <p>由 Spring 的 @Scheduled 注解驱动,在应用启动后自动开始周期性执行。
     * 检测到 OB 不可用时递增失败计数器,达到阈值后触发 failover。</p>
     *
     * <p>注意:如果 autoFailover=false,即使连续失败也不会自动切换,
     * 仅记录警告日志。</p>
     */
    @Scheduled(fixedDelay = 10000)
    public void checkHealth() {
        boolean obOk = checkOceanbase();
        if (!obOk) {
            consecutiveFailures++;
            log.warn("OceanBase 健康检查失败(第 {}/{} 次),当前主库:{}",
                    consecutiveFailures, maxFailures,
                    router.isMasterActive() ? "MySQL" : "OceanBase");
            if (consecutiveFailures >= maxFailures && properties.isAutoFailover()) {
                router.failover();
            }
        } else {
            if (consecutiveFailures > 0) {
                log.info("OceanBase 健康检查恢复,连续失败计数器已重置");
                consecutiveFailures = 0;
                if (properties.isAutoFailover()) {
                    router.recover();
                }
            }
        }
    }

    /**
     * 通过建立 JDBC 连接执行 SELECT 1 来检测 OceanBase 连通性
     *
     * <p>这是最直接的检测方式,比 ping 命令更可靠。
     * 每次检测都会建立真实的数据库连接,确保检测结果准确。</p>
     *
     * @return true = OB 可以正常读写,false = 连接失败
     */
    private boolean checkOceanbase() {
        try (Connection conn = oceanbaseDs.getConnection();
             Statement stmt = conn.createStatement()) {
            stmt.execute("SELECT 1");
            return true;
        } catch (Exception e) {
            return false;
        }
    }

    /**
     * 获取 OceanBase 健康状态(供其他组件快速判断)
     *
     * @return true = 上次检测 OB 正常,false = 上次检测 OB 不可用
     */
    public boolean isOceanbaseHealthy() {
        return oceanbaseHealthy;
    }
}

2.14 log/entity/DualWriteLog.java — 双写日志实体

package com.ruoyi.system.dualwrite.log.entity;

import com.mybatisflex.annotation.Id;
import com.mybatisflex.annotation.KeyType;
import com.mybatisflex.annotation.Table;

import java.util.Date;

/**
 * 双写日志实体类
 *
 * <p>映射数据库表 <b>dual_write_log</b>,记录每次双写同步的详细执行信息。
 * 该表存储在 MySQL 中,双写拦截器会自动排除该表,
 * 因此不会递归同步到 OceanBase。</p>
 *
 * <p>日志用途:</p>
 * <ul>
 *   <li>监控双写成功率</li>
 *   <li>排查双写失败原因</li>
 *   <li>补偿定时任务扫描 FAILED 记录进行重试</li>
 *   <li>统计各表的双写频次</li>
 * </ul>
 *
 * @author ruoyi
 */
@Table("dual_write_log")
public class DualWriteLog {

    /**
     * 日志 ID(数据库自增主键)
     * <p>不使用雪花算法,因为该表不需要双写到 OB,保持简单的自增即可。</p>
     */
    @Id(keyType = KeyType.Auto)
    private Long id;

    /** 被操作的表名,如 sys_user、sys_dept 等 */
    private String tableName;

    /** 操作类型:INSERT、UPDATE、DELETE 之一 */
    private String operation;

    /** 目标数据源,固定为 "OceanBase" */
    private String targetDs;

    /**
     * 同步状态
     * <ul>
     *   <li><b>PENDING</b> — 等待执行(暂未使用,预留)</li>
     *   <li><b>SUCCESS</b> — 同步成功</li>
     *   <li><b>FAILED</b> — 同步失败,错误信息在 errorMsg 字段</li>
     * </ul>
     */
    private String status;

    /** 错误信息,status=FAILED 时记录异常堆栈摘要 */
    private String errorMsg;

    /** 记录创建时间,由数据库 default current_timestamp 自动填充 */
    private Date createTime;

    public Long getId() { return id; }
    public void setId(Long id) { this.id = id; }
    public String getTableName() { return tableName; }
    public void setTableName(String tableName) { this.tableName = tableName; }
    public String getOperation() { return operation; }
    public void setOperation(String operation) { this.operation = operation; }
    public String getTargetDs() { return targetDs; }
    public void setTargetDs(String targetDs) { this.targetDs = targetDs; }
    public String getStatus() { return status; }
    public void setStatus(String status) { this.status = status; }
    public String getErrorMsg() { return errorMsg; }
    public void setErrorMsg(String errorMsg) { this.errorMsg = errorMsg; }
    public Date getCreateTime() { return createTime; }
    public void setCreateTime(Date createTime) { this.createTime = createTime; }
}

2.15 log/service/DualWriteLogService.java — 日志 Service(JdbcTemplate 实现)

package com.ruoyi.system.dualwrite.log.service;

import com.ruoyi.system.dualwrite.log.entity.DualWriteLog;
import org.springframework.jdbc.core.BeanPropertyRowMapper;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;

import java.util.Date;
import java.util.List;

@Service
public class DualWriteLogService {

    private final JdbcTemplate jdbcTemplate;

    public DualWriteLogService(JdbcTemplate jdbcTemplate) {
        this.jdbcTemplate = jdbcTemplate;
    }

    public void saveLog(String tableName, String operation, String targetDs,
                         String status, String errorMsg) {
        try {
            jdbcTemplate.update(
                "INSERT INTO dual_write_log(table_name, operation, target_ds, status, error_msg, create_time) VALUES(?,?,?,?,?,?)",
                tableName, operation, targetDs, status, errorMsg, new Date());
        } catch (Exception e) {
            // ignore
        }
    }

    public List<DualWriteLog> getFailedLogs(int limit) {
        return jdbcTemplate.query(
            "SELECT * FROM dual_write_log WHERE status = 'FAILED' ORDER BY id ASC LIMIT ?",
            new BeanPropertyRowMapper<>(DualWriteLog.class), limit);
    }

    public void updateStatus(Long id, String status, String errorMsg) {
        jdbcTemplate.update(
            "UPDATE dual_write_log SET status = ?, error_msg = ? WHERE id = ?",
            status, errorMsg, id);
    }
}

2.16 job/DualWriteLogCompensationJob.java — 补偿定时任务

package com.ruoyi.system.dualwrite.job;

import com.ruoyi.system.dualwrite.log.entity.DualWriteLog;
import com.ruoyi.system.dualwrite.log.service.DualWriteLogService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;

import java.util.List;

/**
 * 双写失败补偿定时任务
 *
 * <p>每 5 秒扫描 dual_write_log 表中 status = 'FAILED' 的记录,
 * 对失败的双写操作进行补偿重试。</p>
 *
 * <p>当前实现:</p>
 * <ul>
 *   <li>仅输出警告日志,打点记录需要补偿的失败记录</li>
 *   <li>TODO: 后续需要扩展为从日志中解析出 SQL 和参数并重新执行</li>
 * </ul>
 *
 * <p>补偿任务只在 dual-write.enabled=true 时启用,
 * 避免关闭双写后仍然扫描日志表。</p>
 *
 * @author ruoyi
 */
@Component
@ConditionalOnProperty(prefix = "dual-write", name = "enabled", havingValue = "true", matchIfMissing = false)
public class DualWriteLogCompensationJob {

    private static final Logger log = LoggerFactory.getLogger(DualWriteLogCompensationJob.class);

    private final DualWriteLogService logService;

    public DualWriteLogCompensationJob(DualWriteLogService logService) {
        this.logService = logService;
    }

    /**
     * 定时补偿:每 5 秒执行一次
     *
     * <p>由 Spring 的 @Scheduled 注解驱动,应用启动后自动开始周期性执行。
     * 每次最多处理 10 条失败记录,避免一次补偿过多影响系统性能。</p>
     */
    @Scheduled(fixedDelay = 5000)
    public void compensate() {
        List<DualWriteLog> failedLogs = logService.getFailedLogs(10);
        for (DualWriteLog dwl : failedLogs) {
            log.warn("需要补偿的双写记录 ID={},表={},操作={},错误={}",
                    dwl.getId(), dwl.getTableName(),
                    dwl.getOperation(), dwl.getErrorMsg());
            // TODO: 从 dual_write_log 解析出原始 SQL 和参数,重新执行
        }
    }
}

2.17 SQL: resources/mapper/dual_write_log.sql

DROP TABLE IF EXISTS dual_write_log;
CREATE TABLE dual_write_log (
    id            BIGINT(20)    NOT NULL AUTO_INCREMENT  COMMENT '日志 ID(自增主键)',
    table_name    VARCHAR(100)  NOT NULL                 COMMENT '被操作的表名',
    operation     VARCHAR(20)   NOT NULL                 COMMENT 'INSERT / UPDATE / DELETE',
    target_ds     VARCHAR(50)   DEFAULT 'OceanBase'      COMMENT '目标数据源',
    status        VARCHAR(20)   DEFAULT 'PENDING'        COMMENT 'SUCCESS / FAILED',
    error_msg     VARCHAR(2000) DEFAULT NULL             COMMENT '错误信息',
    create_time   DATETIME      DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
    PRIMARY KEY (id)
) ENGINE=INNODB DEFAULT CHARSET=utf8mb4 COMMENT='双写同步日志表';

三、配置文件

3.1 bootstrap.yml(ruoyi-modules/ruoyi-system)

server:
  port: 9201
spring:
  application:
    name: ruoyi-system
  profiles:
    active: dev
  cloud:
    nacos:
      discovery:
        server-addr: 192.168.2.31:8848
      config:
        server-addr: 192.168.2.31:8848
        namespace: 9b4bfe13-cc0a-424b-9ed9-c834df627a35
  config:
    file-extension: yml
    import:
      - nacos:application-${spring.profiles.active}.${spring.config.file-extension}
      - nacos:${spring.application.name}-${spring.profiles.active}.${spring.config.file-extension}
mybatis-flex:
  type-aliases-package: com.ruoyi.system.domain,com.ruoyi.system.api.domain

3.2 ruoyi-system-dev.yml(Nacos 配置中心)

# spring配置
spring:
  data:
    redis:
      host: 192.168.2.31
      port: 6379
      password: 
  datasource:
    url: jdbc:mysql://192.168.2.31:3306/ry?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&useSSL=true&serverTimezone=GMT%2B8
    username: root
    password: '123'
    driver-class-name: com.mysql.cj.jdbc.Driver

oceanbase:
  datasource:
    url: jdbc:mysql://192.168.2.32:2881/ry?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Shanghai
    username: root@sys
    password: '123'
    driver-class-name: com.mysql.cj.jdbc.Driver
    hikari:
      maximum-pool-size: 10
      minimum-idle: 2

dual-write:
  enabled: true
  active: MASTER
  sync-mode: ASYNC
  auto-failover: true

3.3 application-dev.yml(Nacos 配置中心)

spring:
  autoconfigure:
    exclude: com.alibaba.druid.spring.boot3.autoconfigure.DruidDataSourceAutoConfigure
management:
  endpoints:
    web:
      exposure:
        include: '*'

3.4 RuoYiSystemApplication.java(已加 @EnableScheduling

@EnableCustomConfig
@EnableRyFeignClients
@EnableScheduling
@SpringBootApplication
public class RuoYiSystemApplication {
    public static void main(String[] args) {
        SpringApplication.run(RuoYiSystemApplication.class, args);
    }
}

四、Maven 依赖

依赖 GroupId ArtifactId 版本 定义位置 用途
MyBatis-Flex Starter (SB4) com.mybatis-flex mybatis-flex-spring-boot4-starter 1.11.5 根 pom → ruoyi-common-core/pom.xml ORM 框架
MyBatis-Flex Core com.mybatis-flex mybatis-flex-core 1.11.5 根 pom → ruoyi-common-datasource/pom.xml @UseDataSource 注解
Spring JDBC Starter org.springframework.boot spring-boot-starter-jdbc (parent) ruoyi-common-core/pom.xml JdbcTemplate
MySQL Connector/J com.mysql mysql-connector-j (parent) ruoyi-modules-system/pom.xml 驱动
HikariCP com.zaxxer HikariCP (传递) 连接池
Spring Boot Actuator org.springframework.boot spring-boot-starter-actuator (parent) ruoyi-modules-system/pom.xml 健康检查端点

五、遇到的问题及解决方案

5.1 MyBatis-Flex 迁移

  • 全部 Mapper 接口替换为 FlexMapper,实体 @Table 替换 @TableName,分页拦截器替换

5.2 Spring Boot 4.0.6 包路径变更

  • Health/HealthIndicatororg.springframework.boot.health.contributor

5.3 Nacos 旧 Baomidou 配置覆盖

  • Nacos 残留 spring.datasource.dynamic(旧 Druid 多数据源),覆盖了 spring.datasource.url

5.4 Bean 名称冲突

  • 线程池 Bean 名 dualWriteExecutordualWriteTaskExecutor

5.5 循环依赖

  • DualWriteLogService 改用 JdbcTemplate 而非 MyBatis-Flex Mapper

5.6 @Async 名称不匹配

  • DualWriteExecutor.executeAsync()@Async 改为 "dualWriteTaskExecutor"

5.7 MySQL 数据源未被创建(关键 Bug)

  • oceanbaseDs 是唯一 DataSource Bean,Spring Boot 的 @ConditionalOnMissingBean 跳过 MySQL 自动配置
  • 修复:在 DualWriteDataSourceConfig@Primary @Bean(name = "mysqlDs")

5.8 dual_write_log 表不存在

  • 需手动在 MySQL 建表

5.9 SQL 日志不显示参数

  • ms.getBoundSql(null)ms.getBoundSql(parameter)

5.10 Admin 角色无权限

  • sys_role_menu 中 role_id=1 缺少 84 条菜单权限,从 role_id=2 复制补全

六、测试验证

# 登录获取 token
$login = (New-Object System.Net.WebClient).UploadString(
    "http://localhost:9200/login","POST",
    '{"username":"admin","password":"admin123"}')
$token = ($login | ConvertFrom-Json).data.access_token

# 创建配置
$ts = [DateTimeOffset]::Now.ToUnixTimeMilliseconds()
$wc = New-Object System.Net.WebClient
$wc.Headers.Add("Content-Type","application/json")
$wc.Headers.Add("Authorization","Bearer $token")
$wc.UploadString("http://localhost:9201/config","POST",
    "{`"configName`":`"n$ts`",`"configKey`":`"k.$ts`",`"configValue`":`"v$ts`",`"configType`":`"Y`"}")

# 检查 MySQL
$conn = [System.Data.Common.DbProviderFactories]::GetFactory(
    "MySql.Data.MySqlClient").CreateConnection()
$conn.ConnectionString = "server=192.168.2.31;port=3306;database=ry;uid=root;pwd=123"
$conn.Open()
$cmd = $conn.CreateCommand()
$cmd.CommandText = "SELECT config_value FROM sys_config WHERE config_key='k.$ts'"
"MySQL: $($cmd.ExecuteScalar())"
$conn.Close()

# 检查 OceanBase
$conn = [System.Data.Common.DbProviderFactories]::GetFactory(
    "MySql.Data.MySqlClient").CreateConnection()
$conn.ConnectionString = "server=192.168.2.32;port=2881;database=ry;uid=root@sys;pwd=123"
$conn.Open()
$cmd = $conn.CreateCommand()
$cmd.CommandText = "SELECT config_value FROM sys_config WHERE config_key='k.$ts'"
"OB: $($cmd.ExecuteScalar())"
$conn.Close()

# 检查双写日志
$conn = [System.Data.Common.DbProviderFactories]::GetFactory(
    "MySql.Data.MySqlClient").CreateConnection()
$conn.ConnectionString = "server=192.168.2.31;port=3306;database=ry;uid=root;pwd=123"
$conn.Open()
$cmd = $conn.CreateCommand()
$cmd.CommandText = "SELECT id,table_name,operation,status FROM dual_write_log ORDER BY id DESC LIMIT 5"
$reader = $cmd.ExecuteReader()
while ($reader.Read()) {
    "$($reader[0])|$($reader[1])|$($reader[2])|$($reader[3])"
}
$conn.Close()

七、环境信息

组件 地址 说明
MySQL 192.168.2.31:3306 主库,ry 库,root:123
OceanBase 192.168.2.32:2881 从库,ry 库,root@sys:123
Redis 192.168.2.31:6379 缓存
Nacos 192.168.2.31:8848 配置+注册中心,dev 命名空间
JDK D:\install\jdk-21.0.3 JDK 21.0.3
Maven D:\java\apache-maven-3.9.11
MyBatis-Flex v1.11.5 spring-boot4-starter
Spring Boot 4.0.6 Healthorg.springframework.boot.health.contributor
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐