Spring Batch从异源数据库读取时的事务控制方案咨询
Spring Batch跨数据源事务配置(块完成后开启源库事务)
1. 配置多数据源与独立事务管理器
因为源数据库和JobRepository数据库相互独立,需要分别配置两套数据源和对应的事务管理器:
- 一套用于JobRepository的作业状态管理
- 一套用于源数据的读写操作
示例配置代码:
@Configuration public class DataSourceConfig { // JobRepository数据源 @Bean @Primary @ConfigurationProperties(prefix = "spring.datasource.job") public DataSource jobRepositoryDataSource() { return DataSourceBuilder.create().build(); } // 源业务数据库数据源 @Bean @ConfigurationProperties(prefix = "spring.datasource.source") public DataSource sourceDataSource() { return DataSourceBuilder.create().build(); } // JobRepository事务管理器 @Bean public PlatformTransactionManager jobTransactionManager(@Qualifier("jobRepositoryDataSource") DataSource dataSource) { return new DataSourceTransactionManager(dataSource); } // 源数据库事务管理器 @Bean public PlatformTransactionManager sourceTransactionManager(@Qualifier("sourceDataSource") DataSource dataSource) { return new DataSourceTransactionManager(dataSource); } }
接着配置JobRepository绑定专属事务管理器,避免和源库事务混淆:
@Configuration @EnableBatchProcessing public class BatchConfig extends DefaultBatchConfigurer { private final PlatformTransactionManager jobTransactionManager; private final DataSource jobRepositoryDataSource; public BatchConfig(@Qualifier("jobTransactionManager") PlatformTransactionManager jobTransactionManager, @Qualifier("jobRepositoryDataSource") DataSource jobRepositoryDataSource) { this.jobTransactionManager = jobTransactionManager; this.jobRepositoryDataSource = jobRepositoryDataSource; } @Override protected JobRepository createJobRepository() throws Exception { JobRepositoryFactoryBean factory = new JobRepositoryFactoryBean(); factory.setDataSource(jobRepositoryDataSource); factory.setTransactionManager(jobTransactionManager); factory.setIsolationLevelForCreate("ISOLATION_DEFAULT"); factory.setTablePrefix("BATCH_"); return factory.getObject(); } }
2. 自定义非事务型ItemReader
Spring Batch默认的JDBC Reader会在读取阶段开启源库事务,要实现块完成后再开启源库事务,需要关闭Reader的事务感知能力,让读取操作以非事务方式执行:
@Bean public ItemReader<YourEntity> sourceItemReader(@Qualifier("sourceDataSource") DataSource dataSource) { JdbcCursorItemReader<YourEntity> reader = new JdbcCursorItemReader<>(); reader.setDataSource(dataSource); reader.setSql("SELECT id, field1, field2 FROM your_table WHERE processed = false"); reader.setRowMapper((rs, rowNum) -> { YourEntity entity = new YourEntity(); entity.setId(rs.getLong("id")); entity.setField1(rs.getString("field1")); entity.setField2(rs.getString("field2")); return entity; }); // 核心配置:关闭事务感知,读取操作不绑定任何事务 reader.setTransactionAware(false); return reader; }
3. 块完成后手动触发源库事务
在块处理完成(Kafka写入成功)后,通过ChunkListener结合TransactionTemplate手动开启源库事务,执行后续业务操作(比如标记数据已处理):
@Bean public Step dataToKafkaStep(ItemReader<YourEntity> sourceItemReader, ItemWriter<YourEntity> kafkaItemWriter, PlatformTransactionManager sourceTransactionManager, JdbcTemplate sourceJdbcTemplate) { return stepBuilderFactory.get("dataToKafkaStep") .<YourEntity, YourEntity>chunk(100) // 按需调整块大小 .reader(sourceItemReader) .writer(kafkaItemWriter) .listener(new ChunkListener() { private final TransactionTemplate transactionTemplate = new TransactionTemplate(sourceTransactionManager); @Override public void afterChunk(ChunkContext chunkContext) { // 从执行上下文获取当前块处理的所有数据 List<YourEntity> processedItems = (List<YourEntity>) chunkContext.getStepContext() .getStepExecution() .getExecutionContext() .get("processedItems"); // 手动开启源库事务,执行状态更新 transactionTemplate.execute(status -> { sourceJdbcTemplate.update( "UPDATE your_table SET processed = true WHERE id IN (:ids)", new MapSqlParameterSource("ids", processedItems.stream().map(YourEntity::getId).collect(Collectors.toList())) ); return null; }); } }) .build(); }
同时修改Kafka Writer,将处理完成的数据存入执行上下文:
@Bean public ItemWriter<YourEntity> kafkaItemWriter(KafkaTemplate<String, YourEntity> kafkaTemplate) { return items -> { for (YourEntity item : items) { kafkaTemplate.send("your-topic", item.getId().toString(), item); } // 将当前块数据存入执行上下文,供后续Listener使用 StepExecution stepExecution = StepSynchronizationManager.getContext().getStepExecution(); stepExecution.getExecutionContext().put("processedItems", items); }; }
核心逻辑说明
- 源库读取阶段不开启事务,避免长时间占用数据库连接
- JobRepository的事务仅负责作业、步骤的状态管理,与源库事务完全隔离
- 只有当整个块的Kafka写入操作成功后,才会开启源库事务执行状态更新,保证数据一致性
内容的提问来源于stack exchange,提问作者user1409534
相关产品推荐
相关产品推荐

