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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 11:52:46