Spring Batch父子步骤依赖探索:实现T/K类型步骤间数据流交互
基于Spring Batch实现父子步骤依赖的作业流程方案
完全可以借助Spring Batch本身的特性实现你需要的父子步骤数据传递+结果返回流程,以下是具体实现方案:
核心思路
Spring Batch通过JobExecutionContext实现跨步骤的数据共享——父步骤将T类型数据存入全局作业上下文,子步骤读取该数据执行逻辑,完成后再将K类型结果存入上下文,供后续处理器获取。同时利用StepBuilder的监听器机制实现writer后的post-step处理。
具体实现步骤
1. 自定义步骤监听器处理数据传递
通过StepExecutionListener在步骤前后完成数据的读取和存储:
public class DataTransferListener<T, K> implements StepExecutionListener { private T inputData; private K outputData; @Override public void beforeStep(StepExecution stepExecution) { // 从作业上下文读取父步骤传入的T类型数据 JobExecution jobExecution = stepExecution.getJobExecution(); this.inputData = (T) jobExecution.getExecutionContext().get("parentInputData"); } @Override public ExitStatus afterStep(StepExecution stepExecution) { // 将子步骤生成的K类型数据写入作业上下文 JobExecution jobExecution = stepExecution.getJobExecution(); jobExecution.getExecutionContext().put("childOutputData", outputData); return ExitStatus.COMPLETED; } // getter/setter public void setOutputData(K outputData) { this.outputData = outputData; } public T getInputData() { return inputData; } }
2. 配置父子步骤及作业
@Configuration public class ParentChildBatchConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; // 父步骤:生成T类型数据并写入作业上下文 @Bean public Step parentStep() { return stepBuilderFactory.get("parentStep") .tasklet((contribution, chunkContext) -> { // 业务逻辑生成T类型数据 T tData = generateTData(); // 存入全局作业上下文 chunkContext.getStepContext().getStepExecution().getJobExecution() .getExecutionContext().put("parentInputData", tData); return RepeatStatus.FINISHED; }) .build(); } // 子步骤:读取T数据,执行逻辑后生成K类型结果 @Bean public Step childStep(DataTransferListener<T, K> dataTransferListener) { return stepBuilderFactory.get("childStep") .<InputItem, OutputItem>chunk(10) .reader(childItemReader(dataTransferListener)) .processor(childItemProcessor()) .writer(childItemWriter(dataTransferListener)) .listener(dataTransferListener) // 绑定数据传递监听器 .build(); } // 子步骤阅读器:基于父步骤传入的T数据生成读取源 @Bean public ItemReader<InputItem> childItemReader(DataTransferListener<T, K> listener) { return () -> { T parentData = listener.getInputData(); // 根据parentData生成输入项,返回null表示读取结束 return generateInputItem(parentData); }; } // 子步骤处理器:业务数据转换 @Bean public ItemProcessor<InputItem, OutputItem> childItemProcessor() { return inputItem -> { // 业务处理逻辑 return processInputItem(inputItem); }; } // 子步骤Writer:执行写入逻辑并生成K类型结果 @Bean public ItemWriter<OutputItem> childItemWriter(DataTransferListener<T, K> listener) { return items -> { // 写入数据库/文件等操作 writeItems(items); // 生成K类型结果并传递给监听器 K kData = generateKData(items); listener.setOutputData(kData); }; } // 结果处理步骤:读取子步骤返回的K类型数据进行后续处理 @Bean public Step resultProcessorStep() { return stepBuilderFactory.get("resultProcessorStep") .tasklet((contribution, chunkContext) -> { // 从作业上下文读取K类型数据 K kData = (K) chunkContext.getStepContext().getStepExecution().getJobExecution() .getExecutionContext().get("childOutputData"); // 处理K数据的业务逻辑 processKData(kData); return RepeatStatus.FINISHED; }) .build(); } // 组装作业:按顺序执行父步骤→子步骤→结果处理步骤 @Bean public Job parentChildJob() { return jobBuilderFactory.get("parentChildJob") .start(parentStep()) .next(childStep(dataTransferListener())) .next(resultProcessorStep()) .build(); } @Bean public DataTransferListener<T, K> dataTransferListener() { return new DataTransferListener<>(); } // 以下为业务逻辑占位方法,需根据实际场景实现 private T generateTData() { /* ... */ } private InputItem generateInputItem(T parentData) { /* ... */ } private OutputItem processInputItem(InputItem inputItem) { /* ... */ } private void writeItems(List<OutputItem> items) { /* ... */ } private K generateKData(List<OutputItem> items) { /* ... */ } private void processKData(K kData) { /* ... */ } }
3. 关键注意事项
- 序列化要求:存入
JobExecutionContext的T、K类型对象必须实现Serializable接口,否则会抛出序列化异常。 - 步骤顺序:作业中必须保证父步骤先执行,子步骤和结果处理步骤按依赖顺序排列。
- Post-Step处理:如果需要在writer执行完成后立即做局部处理,可直接在
ItemWriter的write方法末尾实现;若需步骤级别的全局处理,使用StepExecutionListener的afterStep方法更合适。
内容的提问来源于stack exchange,提问作者ifrit_prog
相关产品推荐
相关产品推荐

