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

微服务架构下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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 19:04:53