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

Spring Batch处理中如何实现分布式事务?跨多库插入需求求助

跨多数据源多表分布式插入实现方案(基于Spring Batch)

核心思路

利用已配置好的多数据源,为每个数据库的目标表实现独立的ItemWriter,再通过Spring Batch的CompositeItemWriter将多个Writer组合,让Step执行时依次完成跨库插入操作。


1. 为各数据源绑定专属操作模板

假设已通过配置类定义了db1DataSource和db2DataSource两个数据源实例,为每个数据源创建独立的JdbcTemplate(如果用JPA则创建EntityManager):

@Configuration
public class DataSourceConfig {
    @Bean(name = "db1JdbcTemplate")
    public JdbcTemplate db1JdbcTemplate(@Qualifier("db1DataSource") DataSource dataSource) {
        return new JdbcTemplate(dataSource);
    }

    @Bean(name = "db2JdbcTemplate")
    public JdbcTemplate db2JdbcTemplate(@Qualifier("db2DataSource") DataSource dataSource) {
        return new JdbcTemplate(dataSource);
    }
}

2. 实现针对不同表的ItemWriter

分别编写插入db1.tbl_customer和db2.tbl_order的Writer逻辑:

针对db1.tbl_customer的Writer

@Component
public class CustomerItemWriter implements ItemWriter<Customer> {
    private final JdbcTemplate db1JdbcTemplate;

    public CustomerItemWriter(@Qualifier("db1JdbcTemplate") JdbcTemplate db1JdbcTemplate) {
        this.db1JdbcTemplate = db1JdbcTemplate;
    }

    @Override
    public void write(List<? extends Customer> items) throws Exception {
        String insertSql = "INSERT INTO tbl_customer (id, name, email) VALUES (?, ?, ?)";
        db1JdbcTemplate.batchUpdate(insertSql, new BatchPreparedStatementSetter() {
            @Override
            public void setValues(PreparedStatement ps, int i) throws SQLException {
                Customer customer = items.get(i);
                ps.setLong(1, customer.getId());
                ps.setString(2, customer.getName());
                ps.setString(3, customer.getEmail());
            }

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

针对db2.tbl_order的Writer

@Component
public class OrderItemWriter implements ItemWriter<Order> {
    private final JdbcTemplate db2JdbcTemplate;

    public OrderItemWriter(@Qualifier("db2JdbcTemplate") JdbcTemplate db2JdbcTemplate) {
        this.db2JdbcTemplate = db2JdbcTemplate;
    }

    @Override
    public void write(List<? extends Order> items) throws Exception {
        String insertSql = "INSERT INTO tbl_order (order_id, customer_id, amount) VALUES (?, ?, ?)";
        db2JdbcTemplate.batchUpdate(insertSql, new BatchPreparedStatementSetter() {
            @Override
            public void setValues(PreparedStatement ps, int i) throws SQLException {
                Order order = items.get(i);
                ps.setLong(1, order.getOrderId());
                ps.setLong(2, order.getCustomerId());
                ps.setBigDecimal(3, order.getAmount());
            }

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

3. 配置复合ItemWriter并绑定到Step

假设你的ItemProcessor输出包含Customer和Order的复合对象(比如ProcessResult),通过包装Writer拆分数据后,用CompositeItemWriter组合执行:

定义复合结果对象(按需)

public class ProcessResult {
    private Customer customer;
    private Order order;
    // getter、setter省略
}

配置CompositeItemWriter与Step

@Configuration
@EnableBatchProcessing
public class BatchConfig {
    private final JobBuilderFactory jobBuilderFactory;
    private final StepBuilderFactory stepBuilderFactory;
    private final CustomerItemWriter customerItemWriter;
    private final OrderItemWriter orderItemWriter;

    public BatchConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory,
                       CustomerItemWriter customerItemWriter, OrderItemWriter orderItemWriter) {
        this.jobBuilderFactory = jobBuilderFactory;
        this.stepBuilderFactory = stepBuilderFactory;
        this.customerItemWriter = customerItemWriter;
        this.orderItemWriter = orderItemWriter;
    }

    @Bean
    public ItemWriter<ProcessResult> compositeItemWriter() {
        CompositeItemWriter<ProcessResult> compositeWriter = new CompositeItemWriter<>();

        // 拆分ProcessResult为Customer列表,传给CustomerWriter
        ItemWriter<ProcessResult> customerWriterWrapper = items -> {
            List<Customer> customers = items.stream()
                    .map(ProcessResult::getCustomer)
                    .collect(Collectors.toList());
            customerItemWriter.write(customers);
        };

        // 拆分ProcessResult为Order列表,传给OrderWriter
        ItemWriter<ProcessResult> orderWriterWrapper = items -> {
            List<Order> orders = items.stream()
                    .map(ProcessResult::getOrder)
                    .collect(Collectors.toList());
            orderItemWriter.write(orders);
        };

        compositeWriter.setDelegates(Arrays.asList(customerWriterWrapper, orderWriterWrapper));
        return compositeWriter;
    }

    @Bean
    public Step crossDbInsertStep(ItemReader<InputData> itemReader,
                                  ItemProcessor<InputData, ProcessResult> itemProcessor) {
        return stepBuilderFactory.get("crossDbInsertStep")
                .<InputData, ProcessResult>chunk(100) // 批次大小根据业务调整
                .reader(itemReader)
                .processor(itemProcessor)
                .writer(compositeItemWriter())
                .build();
    }

    @Bean
    public Job crossDbInsertJob() {
        return jobBuilderFactory.get("crossDbInsertJob")
                .start(crossDbInsertStep(null, null)) // 实际注入对应的Reader和Processor
                .build();
    }
}

4. 事务注意事项

  • 默认Spring Batch Step事务仅绑定单一数据源,跨库操作无法保证强一致性。若需强一致,需引入XA事务框架(如Atomikos、Bitronix),配置XA数据源与JTA事务管理器。
  • 若业务允许最终一致性,可采用本地事务+补偿机制(如失败重试、记录操作日志后续人工核对)。

内容的提问来源于stack exchange,提问作者Amran Hossain

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 01:25:27