Spring Batch并行处理:可否实现分阶段Split并行Flow的作业配置?
可行性结论
你描述的批处理作业配置完全可以基于Spring Batch原生能力落地,不需要二次修改框架核心逻辑。
实现原理说明
Spring Batch的作业编排模型原生支持三类执行节点的自由组合,节点之间默认按配置顺序串行调度,只有Split节点内部的多个执行单元会走异步并行:
- 普通Step节点:按配置顺序串行执行,前一个节点到达终态(成功/失败/跳过)后才会触发下一个节点
- Split并行节点:节点内挂载的多个Flow/Step会交给指定的任务执行器并行调度,框架会自动等待所有内部单元全部执行完成,才会触发Split节点之后的下一个串行节点
- Flow嵌套节点:可以把一组Step封装成独立的Flow单元,既可以单独串行执行,也可以挂载到Split节点内参与并行调度
你需要的执行流直接按顺序拼接节点即可:
- 作业起始节点配置STEP1
- 后续接第一个Split节点,内部挂载FLOW_1、FLOW_2
- Split节点后接串行节点STEP2
- STEP2后接第二个Split节点,内部挂载FLOW_3、FLOW_4
- 后续按照「串行Step -> Split并行块 -> 串行Step」的规则依次拼接,直到配置完STEP3及所有后续节点即可。
针对你提到的「所有FLOW内置Step结构完全一致、仅对接数据源不同」的场景,完全可以把Flow的构建逻辑抽成通用方法,传入不同的数据源参数、Flow标识就能批量生成对应Flow实例,不需要重复编写Step配置代码。
配置注意事项
- Split节点必须显式配置线程池版的
TaskExecutor(比如ThreadPoolTaskExecutor),如果使用默认的同步任务执行器,并行块会退化成串行执行,达不到并行效果 - 并行执行的多个Flow要做好资源隔离,尽量避免多个Flow同时读写同一份共享资源(尤其是同表写入、同文件读写场景),防止出现锁冲突、脏读脏写问题
- 异常传播规则和串行逻辑一致:同一个Split块内任意一个Flow抛出致命异常,整个Split块会判定为执行失败,后续节点默认不会继续执行,符合Spring Batch默认的作业失败中断逻辑
- 如果后续需要调整并行度,只需要增减Split块内挂载的Flow数量、调整线程池核心线程数即可,不需要改动作业整体编排逻辑。
最简配置示例
以下是基于Spring Batch Java Config的核心配置片段,可直接参考扩展:
@Bean public Job parallelFlowJob(JobRepository jobRepository, Step step1, Step step2, Step step3, TaskExecutor batchThreadPoolExecutor) { // 提前通过通用方法构建好不同数据源对应的Flow Flow flow1 = buildDsFlow("flow-1", ds1, jobRepository, transactionManager); Flow flow2 = buildDsFlow("flow-2", ds2, jobRepository, transactionManager); Flow flow3 = buildDsFlow("flow-3", ds3, jobRepository, transactionManager); Flow flow4 = buildDsFlow("flow-4", ds4, jobRepository, transactionManager); return new JobBuilder("parallelFlowJob", jobRepository) .start(step1) // 第一组并行Flow .split(batchThreadPoolExecutor) .add(flow1, flow2) // 串行执行STEP2 .next(step2) // 第二组并行Flow .split(batchThreadPoolExecutor) .add(flow3, flow4) // 后续继续按规则拼接串行Step、并行Split即可 .next(step3) .build(); } /** * 通用Flow构建方法:传入不同数据源即可生成结构一致的处理Flow */ private Flow buildDsFlow(String flowName, DataSource targetDs, JobRepository jobRepository, PlatformTransactionManager txManager) { Step processStep = new StepBuilder(flowName + "-process", jobRepository) .<SourceData, TargetData>chunk(200, txManager) .reader(buildReader(targetDs)) .processor(commonItemProcessor) .writer(buildWriter(targetDs)) .build(); return new FlowBuilder<Flow>(flowName) .start(processStep) // 可继续扩展Flow内其他通用Step,比如数据校验、结果上报等 .build(); }
内容的提问来源于stack exchange,提问作者hieunt89
相关产品推荐
相关产品推荐

