Spring Batch读取完成标识问题:多Parquet文件批处理场景
解决方案:Spring Batch 处理未知总记录数时的无效记录批量写入
方法1:利用StepExecutionListener监听步骤结束事件
这是最直接的方案,无需判断Reader是否读到最后一个Chunk,只需监听整个Step的结束状态,无论总记录数多少,只要Step执行完成(可按需处理正常结束/异常终止场景),就触发无效记录的写入。
- 实现
StepExecutionListener接口(或使用@BeforeStep/@AfterStep注解) - 在监听器的
afterStep方法中,判断Step是否正常完成,再从内存暂存区取出无效记录写入文件 - 注意:内存暂存的无效记录要放在Step级作用域(比如标注
@StepScope的Bean),避免多线程或Step重启时的状态丢失
示例代码片段:
@Component @StepScope public class InvalidRecordsListener implements StepExecutionListener { // 用线程安全集合暂存无效记录,比如ConcurrentLinkedQueue private final Queue<YourRecordType> invalidRecords = new ConcurrentLinkedQueue<>(); // 提供添加无效记录的方法,供Processor调用 public void addInvalidRecord(YourRecordType record) { invalidRecords.add(record); } @Override public ExitStatus afterStep(StepExecution stepExecution) { // 判断Step是否正常完成 if (stepExecution.getStatus() == BatchStatus.COMPLETED) { // 将无效记录写入文件(以ParquetWriter为例) try (ParquetWriter<YourRecordType> writer = ParquetWriter.builder(...)...) { for (YourRecordType record : invalidRecords) { writer.write(record); } } catch (IOException e) { // 处理写入异常,可设置ExitStatus为FAILED return ExitStatus.FAILED.addExitDescription(e.getMessage()); } } return stepExecution.getExitStatus(); } }
Step配置中注册监听器,并让ItemProcessor注入该监听器来添加无效记录:
@Bean public Step parquetProcessingStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, ItemReader<YourRecordType> parquetReader, ItemProcessor<YourRecordType, YourRecordType> validRecordProcessor, ItemWriter<YourRecordType> dbWriter, InvalidRecordsListener invalidRecordsListener) { return new StepBuilder("parquetProcessingStep", jobRepository) .<YourRecordType, YourRecordType>chunk(1000, transactionManager) .reader(parquetReader) .processor(validRecordProcessor) .writer(dbWriter) .listener(invalidRecordsListener) .build(); }
方法2:自定义ItemReader标记最后一个Chunk
如果必须在处理最后一个Chunk时就触发无效记录写入(而非等Step结束),可以自定义包装Reader,跟踪是否读到了最后一条记录:
- 包装原生Parquet ItemReader,在
read()方法中记录是否返回null(表示无更多记录) - 维护布尔状态
isLastChunk,当Reader返回null时标记当前Chunk为最后一个 - 将该状态传递给ItemWriter或Processor,在处理最后一个Chunk时写入无效记录
示例代码片段:
public class TrackingParquetReader implements ItemReader<YourRecordType> { private final ItemReader<YourRecordType> delegate; private boolean isLastChunk = false; private boolean hasReachedEnd = false; public TrackingParquetReader(ItemReader<YourRecordType> delegate) { this.delegate = delegate; } @Override public YourRecordType read() throws Exception { YourRecordType record = delegate.read(); if (record == null && !hasReachedEnd) { hasReachedEnd = true; isLastChunk = true; } else { isLastChunk = false; } return record; } public boolean isLastChunk() { return isLastChunk; } }
Step配置中使用包装Reader,并让ItemWriter注入该Reader判断是否为最后一个Chunk:
@Bean public ItemReader<YourRecordType> trackingParquetReader(ItemReader<YourRecordType> parquetReader) { return new TrackingParquetReader(parquetReader); } @Bean public ItemWriter<YourRecordType> dbWriter(TrackingParquetReader trackingReader, InvalidRecordsHolder invalidRecordsHolder) { return items -> { // 写入有效数据到数据库 writeToDatabase(items); // 判断是否为最后一个Chunk,若是则写入无效记录 if (trackingReader.isLastChunk()) { writeInvalidRecordsToFile(invalidRecordsHolder.getRecords()); } }; }
注意事项
- 内存暂存无效记录时,需关注数据量大小:若无效记录过多可能引发OOM,可改用临时文件暂存(比如每积累一定数量就写入临时文件,最后合并)
- 分布式/分区Step场景下,要确保无效记录暂存是分区隔离的,避免不同分区记录互相干扰
- 异常处理:若Step中途失败,需明确是否保留已暂存的无效记录,可在
beforeStep中初始化暂存容器,在afterStep中根据状态决定写入或清理
内容的提问来源于stack exchange,提问作者DS.Sun
相关产品推荐
相关产品推荐

