Atomikos + MybatisPlus解决多数据源事务一致性问题解决

 更新时间:2024年07月04日 10:10:14   作者:who_am_i__  
在实际项目的开发过程中,我们经常会遇到在同一个项目或微服务中牵涉到使用两个或多个数据源的,本文主要介绍了Atomikos + MybatisPlus解决多数据源事务一致性问题解决,具有一定的参考价值,感兴趣的可以了解一下

多数据源事务

在实际项目的开发过程中,我们经常会遇到在同一个项目或微服务中牵涉到使用两个或多个数据源的,由于每个数据源需要使用不同的事务管理器,而每个事务管理器管理不同的数据源每个数据源只能保证单个数据源内的事物一致性.所以在使用多个数据源的同时带来的常见问题就是多数据源的事务一致性问题.本文通过利用一种常见的分布式事物管理器atomikos来解决此类问题.

Atomikos

Atomikos 是一个Java事务管理解决方案,用于处理分布式事务。它提供了一个可靠和可扩展的事务管理器,可以协调多个资源(如数据库、消息队列等)之间的事务操作,以保证分布式系统的数据一致性与隔离性。

Atomikos 提供了以下主要功能和特点:

  • 分布式事务协调:Atomikos 使用两阶段提交(Two-Phase Commit)协议来确保分布式事务的一致性。它充当协调者角色,与参与者(各个资源管理器)进行协作并决定是否提交或回滚事务。

  • 事务原子性:Atomikos 确保在分布式环境中进行的事务操作以原子方式执行。如果其中任何一个资源的操作失败,Atomikos 将自动回滚整个事务,确保数据的一致性。

  • 分布式数据源和连接池:Atomikos 提供了分布式数据源和连接池,用于管理多个数据库连接和资源。它能够有效地管理连接和提供高性能和可伸缩性。

  • 事务隔离级别:Atomikos 支持不同的事务隔离级别,包括读未提交(Read Uncommitted)、读已提交(Read Committed)、可重复读(Repeatable Read)和串行化(Serializable)。

  • 高可靠性和扩展性:Atomikos 可以在分布式环境中进行集群部署,提供高可靠性和扩展性。多个 Atomikos 事务管理器可以一起工作,以保证负载均衡和容错。

引入相关依赖

 <dependencies>
        <!-- Spring Boot Starter -->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter</artifactId>
        </dependency>

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>

        <!-- atomikos 依赖-->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-jta-atomikos</artifactId>
        </dependency>

        <!-- mybatis plus -->
        <dependency>
            <groupId>com.baomidou</groupId>
            <artifactId>mybatis-plus-boot-starter</artifactId>
            <version>3.4.3.1</version>
        </dependency>

        <!-- mysql 驱动包 -->
        <dependency>
            <groupId>mysql</groupId>
            <artifactId>mysql-connector-java</artifactId>
            <version>8.0.26</version>
        </dependency>

        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
        </dependency>

    </dependencies>

application.yml配置

# 服务端口号
server:
  port:
    8069

# 多数据源配置
spring:
  datasource:
    properties:
      user_db:
        url: jdbc:mysql://192.168.1.18:3306/user_db
        user: root
        password: xxxxxxx
      data_db:
        url: jdbc:mysql://192.168.1.18:3306/data_db
        user: root
        password: xxxxxxx

上述配置中使用spring.datasource.properties配置多个数据源,在配置类中会使用一个Map<String,Map<String,String>>读取多数据源的每个配置项.

多数据源配置类:DataSourceConfiguration

给配置类负责创建多个数据源以及对应的SqlSessionFactory.代码如下:

package personal.gltm.demo.config;

import com.baomidou.mybatisplus.extension.spring.MybatisSqlSessionFactoryBean;
import com.mysql.cj.jdbc.MysqlXADataSource;
import org.apache.ibatis.session.SqlSessionFactory;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.boot.jta.atomikos.AtomikosDataSourceBean;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.DependsOn;
import org.springframework.context.annotation.Primary;
import org.springframework.core.io.support.PathMatchingResourcePatternResolver;

import javax.sql.DataSource;
import java.util.LinkedHashMap;
import java.util.Map;

@Configuration
@ConfigurationProperties("spring.datasource")
public class DataSourceConfiguration {
    private Map<String,Map<String,String>> properties; // spring.datasource.properties

    private Map<String,DataSource> sourceMap; // 用来保存数据源信息

    public Map<String, Map<String,String>> getProperties() {
        return properties;
    }

    public void setProperties(Map<String, Map<String,String>> properties) {
        this.properties = properties;
    }


    /**
     * 添加user_db数据源
     * @return
     */
    @Bean("user_db")
    @DependsOn({"datasourceMap"})
    public DataSource userDatasource(){
        return sourceMap.get("user_db");
    }

    /**
     * 添加user_db  SqlSessionFactory
     * @return
     * @throws Exception
     */
    @Bean("user_db_sql_session_factory")
    @DependsOn("user_db")
    public SqlSessionFactory userDbSqlSessionFactory() throws Exception {
        MybatisSqlSessionFactoryBean bean = new MybatisSqlSessionFactoryBean();
        bean.setDataSource(this.sourceMap.get("user_db"));
        // 设置mapper位置
        bean.setTypeAliasesPackage("personal.gltm.demo.mapper.user");
        // 设置mapper.xml文件的路径
        bean.setMapperLocations(new PathMatchingResourcePatternResolver().getResources("classpath:user/mapper/*.xml"));
        return bean.getObject();
    }


    /**
     * 添加 data_db 数据源
     * @return
     */
    @Bean("data_db")
    @DependsOn({"datasourceMap"})
    public DataSource dataDatasource(){
        return sourceMap.get("data_db");
    }

    /**
     * 添加data_db SqlSessionFactory
     * @return
     * @throws Exception
     */
    @Bean("data_db_sql_session_factory")
    @DependsOn("data_db")
    public SqlSessionFactory dataDbSqlSessionFactory() throws Exception {
        MybatisSqlSessionFactoryBean bean = new MybatisSqlSessionFactoryBean();
        bean.setDataSource(this.sourceMap.get("data_db"));
        // 设置mapper位置
        bean.setTypeAliasesPackage("personal.gltm.demo.mapper.data");

        // 设置mapper.xml文件的路径
        bean.setMapperLocations(new PathMatchingResourcePatternResolver().getResources("classpath:data/mapper/*.xml"));
        return bean.getObject();
    }


    /**
     * 创建XA 数据源并添加到sourceMap中
     * @param dataSourceConfiguration
     * @return
     */
    @Bean("datasourceMap")
    @Primary
    public Map<String, DataSource> datasourceMap(DataSourceConfiguration dataSourceConfiguration){
        Map<String, DataSource> map = new LinkedHashMap<>();
        for(Map.Entry<String,Map<String, String>> entry:dataSourceConfiguration.properties.entrySet()){
            // 读取每个数据源的配置创建datasource
            MysqlXADataSource dataSource = new MysqlXADataSource();
            dataSource.setUrl(entry.getValue().get("url"));
            dataSource.setUser(entry.getValue().get("user"));
            dataSource.setPassword(entry.getValue().get("password"));

            // 创建atomikosDataSource数据源
            AtomikosDataSourceBean atomikosDataSource = new AtomikosDataSourceBean();
            atomikosDataSource.setMaxPoolSize(10);
            atomikosDataSource.setMinPoolSize(5);
            atomikosDataSource.setBeanName(entry.getKey()); // 设置bean的名称
            atomikosDataSource.setXaDataSource(dataSource); // 设置Xa数据源
            atomikosDataSource.setTestQuery("select now()");

            map.put(entry.getKey(), atomikosDataSource);
        }
        this.sourceMap = map;
        return map;
    }

}

上述配置和代码中我们创建了两个mysql数据源分别是user_db和data_db.并且在两个数据源中分别创建了两个不同的数据表:user_db.t_user{id:bigint,user_name:varchar}和data_db. t_data{id:bigint,context:varchar}

多数据源事务管理器配置:AtomikosConfig

package personal.gltm.demo.config;

import com.atomikos.icatch.jta.UserTransactionImp;
import com.atomikos.icatch.jta.UserTransactionManager;
import lombok.SneakyThrows;
import org.mybatis.spring.SqlSessionFactoryBean;
import org.mybatis.spring.annotation.MapperScan;
import org.mybatis.spring.annotation.MapperScans;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.annotation.EnableTransactionManagement;
import org.springframework.transaction.jta.JtaTransactionManager;

import javax.transaction.SystemException;
import javax.transaction.TransactionManager;
import javax.transaction.UserTransaction;

@Configuration
@EnableTransactionManagement
@MapperScans({@MapperScan(basePackages={"personal.gltm.demo.mapper.user"},sqlSessionFactoryRef = "user_db_sql_session_factory"),
        @MapperScan(basePackages = {"personal.gltm.demo.mapper.data"},sqlSessionFactoryRef = "data_db_sql_session_factory")})
public class AtomikosConfig {

    // 用于在应用程序中执行事务的控制操作。
    @Bean(name = "userTransaction")
    @SneakyThrows(Exception.class)
    public UserTransaction userTransaction() throws SystemException {
        final UserTransactionImp userTransactionImp = new UserTransactionImp();
        userTransactionImp.setTransactionTimeout(1000);
        return userTransactionImp;
    }


    // 用于管理和控制分布式事务的整个生命周期。
    @Bean(name = "atomikosTransactionManager")
    @SneakyThrows(Exception.class)
    public TransactionManager atomikosTransactionManager() {
        final UserTransactionManager userTransactionManager = new UserTransactionManager();
        userTransactionManager.setForceShutdown(false);
        return userTransactionManager;
    }


    // JtaTransactionManager 的主要作用是管理和协调分布式事务,它支持使用 JTA 来处理分布式事务,与 JTA 兼容的事务管理器进行交互。
    @Bean(name = "transactionManager")
    @SneakyThrows(Throwable.class)
    public PlatformTransactionManager transactionManager(
            @Qualifier("atomikosTransactionManager") TransactionManager atomikosTransactionManager,
                                                         @Qualifier("userTransaction") UserTransaction userTransaction) throws SystemException {
        return new JtaTransactionManager(userTransaction(), atomikosTransactionManager());
    }

}

上述代码需要使用@MapperScans注解绑定sqlsessionfactory的扫描包. 

验证测试

添加测试类:

package personal.gltm.demo.controller;

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import personal.gltm.demo.model.Data;
import personal.gltm.demo.model.User;
import personal.gltm.demo.service.DataService;
import personal.gltm.demo.service.UserService;

import java.math.BigDecimal;

@RequestMapping("/test")
@RestController
public class TestController {

    @Autowired
    private UserService userService;

    @Autowired
    private DataService dataService;

    @GetMapping("/test/{u-id}/{d-id}")
    @Transactional(rollbackFor = Throwable.class)
    public String Test(@PathVariable("u-id")BigDecimal uId,@PathVariable("d-id") BigDecimal dId){
        User user = new User();
        user.setId(uId);
        user.setUserName("sihong" + Math.floor(Math.random() * 20));
        userService.insert(user);

        Data data = new Data();
        data.setId(dId);
        data.setContext("mine mine mine ..... " + System.nanoTime());
        dataService.insert(data);

        return "success";
    }
}

测试类中,我们添加了一个Test方法,接收两个参数:userId和dataId,然后把这两个参数作为主键分别插入到user_db.t_user和data_db.t_data中.如果主键冲突则其中一个会报错.如果要保证多数据源事物一致性另一个事物也必须回滚.

在浏览器中先调用:http://localhost/test/test/1/1   分别在user_db.t_user 和data_db.t_data中插入主键为1 的数据.再调用http://localhost/test/test/1/2 时会向数据库插入主键分别为1和2的数据,但是t_user中会存在主键冲突.整个事务回滚t_data 中也不会插入数据.

到此这篇关于Atomikos + MybatisPlus解决多数据源事务一致性问题解决的文章就介绍到这了,更多相关Atomikos MybatisPlus多数据源事务一致性内容请搜索脚本之家以前的文章或继续浏览下面的相关文章希望大家以后多多支持脚本之家!

相关文章

  • SpringCloud如何引用xxjob定时任务

    SpringCloud如何引用xxjob定时任务

    Spring Cloud 本身不直接支持 XXL-JOB 这样的定时任务框架,如果你想在 Spring Cloud 应用中集成 XXL-JOB,你需要手动进行配置,本文给大家介绍SpringCloud如何引用xxjob定时任务,感兴趣的朋友一起看看吧
    2024-04-04
  • 在idea中为注释标记作者日期操作

    在idea中为注释标记作者日期操作

    这篇文章主要介绍了在idea中为注释标记作者日期操作,具有很好的参考价值,希望对大家有所帮助。一起跟随小编过来看看吧
    2020-08-08
  • SpringCloud gateway跨域配置的操作

    SpringCloud gateway跨域配置的操作

    这篇文章主要介绍了SpringCloud gateway跨域配置的操作,具有很好的参考价值,希望对大家有所帮助。如有错误或未考虑完全的地方,望不吝赐教
    2021-07-07
  • 详细谈谈Java中long和double的原子性

    详细谈谈Java中long和double的原子性

    原子性是指一个操作或多个操作要么全部执行,且执行的过程不会被任何因素打断,要么就都不执行,下面这篇文章主要给大家介绍了关于Java中long和double原子性的相关资料,需要的朋友可以参考下
    2021-08-08
  • Java设计模式之监听器模式实例详解

    Java设计模式之监听器模式实例详解

    这篇文章主要介绍了Java设计模式之监听器模式,结合实例形式较为详细的分析了java设计模式中监听器模式的概念、原理及相关实现与使用技巧,需要的朋友可以参考下
    2018-02-02
  • spring-retry组件的使用教程

    spring-retry组件的使用教程

    Spring Retry的主要目的是为了提高系统的可靠性和容错性,当方法调用失败时,Spring Retry可以在不影响系统性能的情况下,自动进行重试,从而减少故障对系统的影响,这篇文章主要介绍了spring-retry组件的使用,需要的朋友可以参考下
    2023-06-06
  • 详解Java ScheduledThreadPoolExecutor的踩坑与解决方法

    详解Java ScheduledThreadPoolExecutor的踩坑与解决方法

    最近项目上反馈某个重要的定时任务突然不执行了,很头疼,开发环境和测试环境都没有出现过这个问题。定时任务采用的是ScheduledThreadPoolExecutor,后来一看代码发现踩了一个大坑。本文就来和大家聊聊这次的踩坑记录与解决方法,需要的可以参考一下
    2022-10-10
  • JVM双亲委派模型知识详细总结

    JVM双亲委派模型知识详细总结

    今天带各位小伙伴学习Java虚拟机的相关知识,文中对JVM双亲委派模型作了非常详细的介绍,对正在学习java的小伙伴们有很好的帮助,需要的朋友可以参考下
    2021-05-05
  • Java实现短信验证码和国际短信群发功能的示例

    Java实现短信验证码和国际短信群发功能的示例

    本篇文章主要介绍了Java实现短信验证码和国际短信群发功能的示例,具有一定的参考价值,感兴趣的小伙伴们可以参考一下。
    2017-02-02
  • SpringBoot工程下Lombok的应用教程详解

    SpringBoot工程下Lombok的应用教程详解

    这篇文章主要给大家介绍了关于SpringBoot工程下Lombok应用的相关资料,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧
    2020-11-11

最新评论