Spring Batch:如何让第一个Writer写完所有数据再执行第二个?
核心疑问解答
Step之间可以传递数据吗?
可以。Spring Batch通过JobExecutionContext实现Step间的数据共享,不过要注意共享的数据必须实现Serializable接口,否则会序列化失败。是否必须重新读取数据?
不一定。如果第一个Step已经把需要处理的数据全量收集并存储(内存或临时介质),第二个Step可以直接读取这些预处理好的数据,无需重新从原始数据源读取。但如果数据量极大,内存存储会有性能问题,这时候建议用临时数据库表或文件存储。
解决"ItemReader must be provided"错误&实现需求的具体方案
方案一:内存共享数据+自定义读取Reader(适合小到中等数据量)
思路:第一个Step先执行FirstWriter处理每个Chunk,同时把所有数据收集到内存中,通过JobExecutionContext传递给第二个Step;第二个Step用自定义Reader读取共享的数据,再交给SecondWriter处理。
- 自定义收集型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; } }
- 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; } }
- 自定义读取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; } }
- 配置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,建议用临时介质中转:
- 第一个Step:FirstWriter将处理后的数据写入临时数据库表(或CSV文件)
- 第二个Step:用
JdbcCursorItemReader(或FlatFileItemReader)读取临时存储的数据,交给SecondWriter处理 - 可选:添加第三个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

