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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 23:57:28