微服务架构下Spring Batch Reader用Feign拉取数据的堆溢出及一致性问题
Spring Batch一次性拉取数据堆溢出的解决方案与分页一致性保障
一、解决堆内存溢出的核心方案
你的代码中使用ListItemReader一次性加载所有数据到内存,当数据量较大时必然触发堆溢出。解决思路是替换为分页式自定义ItemReader,分批拉取数据,每次只加载当前Chunk所需的部分数据,避免全量数据驻留内存。
二、分页拉取的数据一致性保障策略
分页拉取时需避免数据新增/修改/删除导致的重复处理、漏处理或版本不一致问题,可根据业务场景选择以下策略:
1. 游标分页(基于唯一递增ID)
适用于只需处理新增数据或按ID顺序处理的场景:
- 以上次处理的最后一条数据ID作为下一页的起始条件,而非传统的offset分页,避免数据新增时offset偏移导致的重复读取。
- 将处理进度(最后处理的ID)持久化到
ExecutionContext,Job重启时可从断点继续,保证不重复、不遗漏。
2. 时间快照分页
适用于需要处理某一时间点全量数据的场景:
- Job启动时记录当前时间作为快照时间,拉取该时间点之前的所有数据。
- 服务端需支持基于时间范围的分页查询,最好配合读锁或快照机制,保证处理的是快照时间点的数据版本,不受后续数据修改影响。
3. 临时表快照
适用于数据源允许的场景:
- Job启动时先将需处理的数据同步到临时表,后续从临时表分页拉取。临时表数据静态隔离,原数据的新增/修改/删除不会影响Job处理。
三、具体实现代码
1. 基于游标分页的自定义ItemReader
@StepScope @Bean(name = "paymentPagingReader") public ItemReader<PaymentDto> paymentPagingReader(@Value("#{jobParameters['pageSize'] ?: 100}") int pageSize) { return new AbstractPaginatedItemReader<PaymentDto>() { private long lastProcessedId = 0; private List<PaymentDto> currentPageData = new ArrayList<>(); private int currentIndex = 0; @Override protected void doReadPage() { // 调用Feign接口,基于最后处理ID拉取分页数据 ServiceResponse<List<PaymentDto>> result = paymentClient.getPaymentListAfterId(lastProcessedId, pageSize); currentPageData = result.getData(); if (!currentPageData.isEmpty()) { // 更新最后处理的ID为当前页最后一条数据的ID lastProcessedId = currentPageData.get(currentPageData.size() - 1).getId(); } currentIndex = 0; } @Override protected PaymentDto doRead() { if (currentIndex >= currentPageData.size()) { doReadPage(); if (currentPageData.isEmpty()) { return null; // 无更多数据,结束读取 } } return currentPageData.get(currentIndex++); } @Override public void open(ExecutionContext executionContext) throws ItemStreamException { super.open(executionContext); // 从ExecutionContext恢复上次中断的处理进度 if (executionContext.containsKey("lastProcessedId")) { this.lastProcessedId = executionContext.getLong("lastProcessedId"); } } @Override public void update(ExecutionContext executionContext) throws ItemStreamException { super.update(executionContext); // 将当前处理进度保存到ExecutionContext executionContext.putLong("lastProcessedId", this.lastProcessedId); } }; }
2. 修改Step配置
替换原ListItemReader为自定义分页Reader:
private Step createStepFirst(ItemWriter<PaymentDto> createWriter, ItemProcessor<PaymentDto, PaymentDto> itemProcessor, ItemReader<PaymentDto> reader) { return stepBuilderFactory.get(StepName.STEPFIRST) .<PaymentDto, PaymentDto>chunk(100) .reader(reader) .processor(itemProcessor) .writer(createWriter) .build(); }
3. 时间快照分页的Reader示例
@StepScope @Bean(name = "paymentSnapshotReader") public ItemReader<PaymentDto> paymentSnapshotReader(@Value("#{jobParameters['snapshotTime']}") Date snapshotTime, @Value("#{jobParameters['pageSize'] ?: 100}") int pageSize) { return new AbstractPaginatedItemReader<PaymentDto>() { private int currentPage = 0; private List<PaymentDto> currentPageData = new ArrayList<>(); private int currentIndex = 0; @Override protected void doReadPage() { // 拉取快照时间之前的分页数据 ServiceResponse<List<PaymentDto>> result = paymentClient.getPaymentListBeforeTime(snapshotTime, currentPage, pageSize); currentPageData = result.getData(); currentPage++; currentIndex = 0; } @Override protected PaymentDto doRead() { if (currentIndex >= currentPageData.size()) { doReadPage(); if (currentPageData.isEmpty()) { return null; } } return currentPageData.get(currentIndex++); } }; }
4. 启动Job时传入快照时间(时间快照场景)
JobParameters jobParameters = new JobParametersBuilder() .addDate("snapshotTime", new Date()) .addLong("pageSize", 100L) .toJobParameters(); jobLauncher.run(paymentJob, jobParameters);
内容的提问来源于stack exchange,提问作者kernel
相关产品推荐
相关产品推荐

