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

Spring Batch:如何让第一个Writer写完所有数据再执行第二个?

Spring Batch 4.3.7:实现先完成所有Chunk的FirstWriter后再执行SecondWriter的方案

核心疑问解答

  • Step之间可以传递数据吗?
    可以。Spring Batch通过JobExecutionContext实现Step间的数据共享,不过要注意共享的数据必须实现Serializable接口,否则会序列化失败。

  • 是否必须重新读取数据?
    不一定。如果第一个Step已经把需要处理的数据全量收集并存储(内存或临时介质),第二个Step可以直接读取这些预处理好的数据,无需重新从原始数据源读取。但如果数据量极大,内存存储会有性能问题,这时候建议用临时数据库表或文件存储。


解决"ItemReader must be provided"错误&实现需求的具体方案

方案一:内存共享数据+自定义读取Reader(适合小到中等数据量)

思路:第一个Step先执行FirstWriter处理每个Chunk,同时把所有数据收集到内存中,通过JobExecutionContext传递给第二个Step;第二个Step用自定义Reader读取共享的数据,再交给SecondWriter处理。

  1. 自定义收集型FirstWriter
public class CollectingFirstWriter<T> implements ItemWriter<T> {
    private final List<T> collectedItems = new ArrayList<>();
    private final ItemWriter<T> originalFirstWriter; // 原来的FirstWriter

    public CollectingFirstWriter(ItemWriter<T> originalFirstWriter) {
        this.originalFirstWriter = originalFirstWriter;
    }

    @Override
    public void write(List<? extends T> items) throws Exception {
        // 先执行原FirstWriter的逻辑
        originalFirstWriter.write(items);
        // 累积所有Chunk的数据
        collectedItems.addAll(items);
    }

    public List<T> getCollectedItems() {
        return collectedItems;
    }
}
  1. Step监听类,将收集的数据写入JobExecutionContext
public class CollectingStepListener implements StepExecutionListener {
    private final CollectingFirstWriter<?> collectingWriter;

    public CollectingStepListener(CollectingFirstWriter<?> collectingWriter) {
        this.collectingWriter = collectingWriter;
    }

    @Override
    public ExitStatus afterStep(StepExecution stepExecution) {
        ExecutionContext jobContext = stepExecution.getJobExecution().getExecutionContext();
        // 将收集的数据存入Job上下文
        jobContext.put("processedItems", collectingWriter.getCollectedItems());
        return ExitStatus.COMPLETED;
    }
}
  1. 自定义读取Job上下文数据的Reader
public class ExecutionContextItemReader<T> implements ItemReader<T> {
    private Iterator<T> itemIterator;

    @BeforeStep
    public void init(StepExecution stepExecution) {
        ExecutionContext jobContext = stepExecution.getJobExecution().getExecutionContext();
        List<T> items = (List<T>) jobContext.get("processedItems");
        itemIterator = items.iterator();
    }

    @Override
    public T read() {
        return itemIterator.hasNext() ? itemIterator.next() : null;
    }
}
  1. 配置Job和Steps
@Configuration
public class BatchJobConfig {
    // 第一个Step:处理数据+收集
    @Bean
    public Step processAndCollectStep(JobRepository jobRepository,
                                      PlatformTransactionManager transactionManager,
                                      ItemReader<YourItemType> originalReader,
                                      ItemProcessor<YourItemType, YourItemType> originalProcessor,
                                      ItemWriter<YourItemType> originalFirstWriter) {
        CollectingFirstWriter<YourItemType> collectingWriter = new CollectingFirstWriter<>(originalFirstWriter);
        CollectingStepListener stepListener = new CollectingStepListener(collectingWriter);

        return new StepBuilder("processAndCollectStep", jobRepository)
                .<YourItemType, YourItemType>chunk(10, transactionManager)
                .reader(originalReader)
                .processor(originalProcessor)
                .writer(collectingWriter)
                .listener(stepListener)
                .build();
    }

    // 第二个Step:读取收集的数据+执行SecondWriter
    @Bean
    public Step finalWriteStep(JobRepository jobRepository,
                               PlatformTransactionManager transactionManager,
                               ItemWriter<YourItemType> secondWriter) {
        return new StepBuilder("finalWriteStep", jobRepository)
                .<YourItemType, YourItemType>chunk(10, transactionManager)
                .reader(new ExecutionContextItemReader<>())
                .writer(secondWriter)
                .build();
    }

    // 组装Job
    @Bean
    public Job sequentialWriterJob(JobRepository jobRepository,
                                   Step processAndCollectStep,
                                   Step finalWriteStep) {
        return new JobBuilder("sequentialWriterJob", jobRepository)
                .start(processAndCollectStep)
                .next(finalWriteStep)
                .build();
    }
}

方案二:临时存储(适合大数据量)

如果数据量太大,内存存储会导致OOM,建议用临时介质中转:

  1. 第一个Step:FirstWriter将处理后的数据写入临时数据库表(或CSV文件)
  2. 第二个Step:用JdbcCursorItemReader(或FlatFileItemReader)读取临时存储的数据,交给SecondWriter处理
  3. 可选:添加第三个Step,清理临时存储的数据

方案三:自定义延迟执行的CompositeWriter(无需拆分Step)

如果不想拆分成两个Step,可以自定义一个CompositeWriter,先执行FirstWriter处理每个Chunk,等所有Chunk处理完成后再一次性执行SecondWriter:

public class DelayedCompositeWriter<T> implements ItemWriter<T>, StepExecutionListener {
    private final ItemWriter<T> firstWriter;
    private final ItemWriter<T> secondWriter;
    private final List<T> allProcessedItems = new ArrayList<>();

    public DelayedCompositeWriter(ItemWriter<T> firstWriter, ItemWriter<T> secondWriter) {
        this.firstWriter = firstWriter;
        this.secondWriter = secondWriter;
    }

    @Override
    public void write(List<? extends T> items) throws Exception {
        // 每个Chunk先执行FirstWriter
        firstWriter.write(items);
        // 累积数据
        allProcessedItems.addAll(items);
    }

    @Override
    public ExitStatus afterStep(StepExecution stepExecution) {
        try {
            // 所有Chunk处理完后,执行SecondWriter
            secondWriter.write(allProcessedItems);
        } catch (Exception e) {
            return ExitStatus.FAILED.addExitDescription(e.getMessage());
        }
        return ExitStatus.COMPLETED;
    }
}

直接用这个Writer替换原来的CompositeItemWriter即可,无需修改Step结构。


关于报错的说明

java.lang.IllegalStateException: ItemReader must be provided是因为Spring Batch的Chunk-oriented Step要求必须配置ItemReader。上面的方案一、二都为第二个Step提供了合法的Reader,方案三无需拆分Step,自然不会触发这个错误。

内容的提问来源于stack exchange,提问作者Czarek

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 03:05:56