You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用JdbcTemplate批量更新多数据源时,出错如何实现数据回滚?

跨数据源批量更新的事务回滚解决方案

你当前代码的核心问题在于Repository层使用了REQUIRES_NEW事务传播属性,导致每个批量操作都会开启独立事务并立即提交,上层Service的事务无法控制这些子事务的回滚。要实现两个数据源的操作同步(要么都成功,要么都回滚),需要通过分布式事务管理来解决。

1. 移除Repository层的独立事务配置

首先删掉CustomerRepository和AddressRepository中的@Transactional(propagation = Propagation.REQUIRES_NEW),让它们的操作纳入上层Service的事务上下文:

修改后的CustomerRepository.java

@Repository
public class CustomerRepository {
    private final JdbcTemplate jdbcTemplate;

    public CustomerRepository(@Qualifier("customerJdbcTemplate") JdbcTemplate jdbcTemplate) {
        this.jdbcTemplate = jdbcTemplate;
    }

    public int[] saveSuccessBatch(List<Customer> customerList) {
        return jdbcTemplate.batchUpdate("insert into customer (age,first_name, last_name) values(?,?,?)",
            new BatchPreparedStatementSetter() {
                @Override
                public void setValues(PreparedStatement ps, int i) throws SQLException {
                    Customer customer = customerList.get(i);
                    ps.setInt(1, customer.getAge());
                    ps.setString(2, customer.getFirstName());
                    ps.setString(3, customer.getLastName());
                }

                @Override
                public int getBatchSize() {
                    return customerList.size();
                }
            });
    }
}

修改后的AddressRepository.java

@Repository
public class AddressRepository {
    private final JdbcTemplate jdbcTemplate;

    public AddressRepository(@Qualifier("addressJdbcTemplate") JdbcTemplate jdbcTemplate) {
        this.jdbcTemplate = jdbcTemplate;
    }

    public int[] saveSuccessBatch(List<Address> addressList) {
        return jdbcTemplate.batchUpdate("insert into address (city,street, zipCode) values(?,?,?)",
            new BatchPreparedStatementSetter() {
                @Override
                public void setValues(PreparedStatement ps, int i) throws SQLException {
                    Address address = addressList.get(i);
                    ps.setString(1, address.getCity());
                    ps.setString(2, address.getStreet());
                    ps.setInt(3, address.getZipCode());
                }

                @Override
                public int getBatchSize() {
                    return addressList.size();
                }
            });
    }
}

2. 配置JTA分布式事务管理器(以Atomikos为例)

跨数据源事务需要依赖JTA(Java Transaction API)实现全局事务管理,这里以Atomikos作为JTA实现:

2.1 添加Maven依赖

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

2.2 配置双数据源与JTA事务管理器

创建配置类,用Atomikos包装两个数据源,并配置全局事务管理器:

@Configuration
public class DataSourceConfig {

    // 客户库数据源
    @Bean(name = "customerDataSource")
    public DataSource customerDataSource(Environment env) {
        AtomikosDataSourceBean ds = new AtomikosDataSourceBean();
        ds.setUniqueResourceName("customerDB");
        ds.setXaDataSourceClassName("com.mysql.cj.jdbc.MysqlXADataSource");
        Properties props = new Properties();
        props.put("url", env.getProperty("spring.datasource.customer.url"));
        props.put("user", env.getProperty("spring.datasource.customer.username"));
        props.put("password", env.getProperty("spring.datasource.customer.password"));
        ds.setXaProperties(props);
        ds.setMinPoolSize(5);
        ds.setMaxPoolSize(20);
        return ds;
    }

    // 地址库数据源
    @Bean(name = "addressDataSource")
    public DataSource addressDataSource(Environment env) {
        AtomikosDataSourceBean ds = new AtomikosDataSourceBean();
        ds.setUniqueResourceName("addressDB");
        ds.setXaDataSourceClassName("com.mysql.cj.jdbc.MysqlXADataSource");
        Properties props = new Properties();
        props.put("url", env.getProperty("spring.datasource.address.url"));
        props.put("user", env.getProperty("spring.datasource.address.username"));
        props.put("password", env.getProperty("spring.datasource.address.password"));
        ds.setXaProperties(props);
        ds.setMinPoolSize(5);
        ds.setMaxPoolSize(20);
        return ds;
    }

    // 对应客户库的JdbcTemplate
    @Bean(name = "customerJdbcTemplate")
    public JdbcTemplate customerJdbcTemplate(@Qualifier("customerDataSource") DataSource dataSource) {
        return new JdbcTemplate(dataSource);
    }

    // 对应地址库的JdbcTemplate
    @Bean(name = "addressJdbcTemplate")
    public JdbcTemplate addressJdbcTemplate(@Qualifier("addressDataSource") DataSource dataSource) {
        return new JdbcTemplate(dataSource);
    }

    // 全局JTA事务管理器
    @Bean
    public PlatformTransactionManager transactionManager() {
        UserTransactionManager userTransactionManager = new UserTransactionManager();
        UserTransaction userTransaction = new UserTransactionImp();
        return new JtaTransactionManager(userTransaction, userTransactionManager);
    }
}

3. 调整Service层事务配置

注意:Spring AOP无法拦截private方法,所以需要将saveData改为public,确保事务能覆盖所有操作:

修改后的Service.java

@Service
public class DataService {
    private final CustomerRepository customerRepository;
    private final AddressRepository addressRepository;

    public DataService(CustomerRepository customerRepository, AddressRepository addressRepository) {
        this.customerRepository = customerRepository;
        this.addressRepository = addressRepository;
    }

    // 全局事务,覆盖所有批量操作
    @Transactional(rollbackFor = Exception.class)
    public void saveDataInBatch() {
        List<List<Customer>> mainCustomerList = getMainCustomerList();
        List<List<Address>> mainAddressList = getMainAddressList();

        for (int index = 0; index < 3; index++) {
            saveData(mainCustomerList, mainAddressList, index);
        }
    }

    // 改为public,让事务能拦截
    public void saveData(List<List<Customer>> mainCustomerList, List<List<Address>> mainAddressList, int index) {
        customerRepository.saveSuccessBatch(mainCustomerList.get(index));
        addressRepository.saveSuccessBatch(mainAddressList.get(index));
    }

    private List<List<Customer>> getMainCustomerList() {
        // 模拟业务数据
        return new ArrayList<>();
    }

    private List<List<Address>> getMainAddressList() {
        // 模拟业务数据
        return new ArrayList<>();
    }
}

关键注意事项

  • 分布式事务的必要性:单数据源事务无法跨多个独立数据库,必须依赖JTA或Seata等分布式事务框架。
  • 事务传播属性:禁止在Repository层使用REQUIRES_NEW,否则会脱离上层事务的控制。
  • 异常回滚:添加rollbackFor = Exception.class确保所有异常都能触发回滚(默认仅RuntimeException会回滚)。
  • private方法事务失效:Spring AOP只能代理public方法,事务注解不能放在private方法上。

内容的提问来源于stack exchange,提问作者Dharmin Patel

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 17:42:36