Spring Boot Batch实现步骤B&C并行执行的配置方案问询
Spring Batch 作业并行化改造:A→并行B/C→D
一、核心配置确认
- Quartz单工作线程无需调整:Quartz的单线程仅负责定时触发作业启动,作业内部的步骤并行由Spring Batch的TaskExecutor处理,只要Quartz能正常触发作业即可。
- 必须配置Spring Batch多线程执行器:通过
TaskExecutor实现B、C步骤的并行运行。
二、当前串行作业配置示例
@Configuration public class SerialBatchConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; // 步骤A:截断目标表 @Bean public Step stepA(DataSource dataSource) { return stepBuilderFactory.get("stepA") .tasklet((contribution, chunkContext) -> { try (Connection conn = dataSource.getConnection()) { conn.createStatement().execute("TRUNCATE TABLE target_table"); } return RepeatStatus.FINISHED; }) .build(); } // 步骤B:Chunk大小2000的读写组件 @Bean public Step stepB(JdbcCursorItemReader<MyEntity> readerB, JdbcBatchItemWriter<MyEntity> writerB) { return stepBuilderFactory.get("stepB") .<MyEntity, MyEntity>chunk(2000) .reader(readerB) .writer(writerB) .build(); } // 步骤C:Chunk大小2000的读写组件 @Bean public Step stepC(JdbcCursorItemReader<MyEntity> readerC, JdbcBatchItemWriter<MyEntity> writerC) { return stepBuilderFactory.get("stepC") .<MyEntity, MyEntity>chunk(2000) .reader(readerC) .writer(writerC) .build(); } // 步骤D:执行PL/SQL存储过程 @Bean public Step stepD(DataSource dataSource) { return stepBuilderFactory.get("stepD") .tasklet((contribution, chunkContext) -> { try (Connection conn = dataSource.getConnection()) { CallableStatement cs = conn.prepareCall("{CALL post_process_procedure()}"); cs.execute(); } return RepeatStatus.FINISHED; }) .build(); } // 串行作业流程:A→B→C→D @Bean public Job serialJob() { return jobBuilderFactory.get("serialJob") .start(stepA(null)) .next(stepB(null, null)) .next(stepC(null, null)) .next(stepD(null)) .build(); } }
三、并行化改造后的作业配置示例
@Configuration public class ParallelBatchConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; // 配置多线程执行器,并发数设为2刚好支持B、C并行 @Bean public TaskExecutor taskExecutor() { SimpleAsyncTaskExecutor taskExecutor = new SimpleAsyncTaskExecutor(); taskExecutor.setConcurrencyLimit(2); return taskExecutor; } // 步骤A:与串行版本完全一致 @Bean public Step stepA(DataSource dataSource) { return stepBuilderFactory.get("stepA") .tasklet((contribution, chunkContext) -> { try (Connection conn = dataSource.getConnection()) { conn.createStatement().execute("TRUNCATE TABLE target_table"); } return RepeatStatus.FINISHED; }) .build(); } // 步骤B:与串行版本完全一致 @Bean public Step stepB(JdbcCursorItemReader<MyEntity> readerB, JdbcBatchItemWriter<MyEntity> writerB) { return stepBuilderFactory.get("stepB") .<MyEntity, MyEntity>chunk(2000) .reader(readerB) .writer(writerB) .build(); } // 步骤C:与串行版本完全一致 @Bean public Step stepC(JdbcCursorItemReader<MyEntity> readerC, JdbcBatchItemWriter<MyEntity> writerC) { return stepBuilderFactory.get("stepC") .<MyEntity, MyEntity>chunk(2000) .reader(readerC) .writer(writerC) .build(); } // 步骤D:与串行版本完全一致 @Bean public Step stepD(DataSource dataSource) { return stepBuilderFactory.get("stepD") .tasklet((contribution, chunkContext) -> { try (Connection conn = dataSource.getConnection()) { CallableStatement cs = conn.prepareCall("{CALL post_process_procedure()}"); cs.execute(); } return RepeatStatus.FINISHED; }) .build(); } // 定义并行Flow:同时执行B、C @Bean public Flow parallelFlow() { return new FlowBuilder<SimpleFlow>("parallelFlow") .split(taskExecutor()) .add(stepB(null, null), stepC(null, null)) .build(); } // 最终并行作业流程:A→并行B/C→D @Bean public Job parallelJob() { return jobBuilderFactory.get("parallelJob") .start(stepA(null)) .next(parallelFlow()) .next(stepD(null)) .end() .build(); } }
四、关键注意事项
- 步骤独立性:确保B、C操作无数据冲突(比如读写不同表、无重叠数据范围),避免并行时出现锁或数据不一致问题。
- Reader线程安全:JdbcCursorItemReader本身非线程安全,但Spring Batch会为并行的每个步骤线程创建独立的Reader实例,无需额外处理。
- 执行器选择:示例用
SimpleAsyncTaskExecutor,生产环境可根据需求替换为ThreadPoolTaskExecutor,更适合稳定的多线程场景。
内容的提问来源于stack exchange,提问作者Lokesh Parihar
相关产品推荐
相关产品推荐

