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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 04:06:08