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
相关产品推荐
相关产品推荐

