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

Spring Batch作业分页读取时页面大量丢失问题求助

问题描述

背景

使用分页方式从数据库读取数据,采用batch size为1000的RepositoryItemReader,基于created_at和status字段过滤数据。读取事务后需将其状态更新为CANCELLED,再发送至Kafka主题供其他服务监听,但目前并非所有事务都能被正常更新。

核心问题

由于更新后的事务状态恰好是RepositoryItemReader的过滤条件,导致作业跳过了近一半的页面,大量数据漏处理。

期望目标

让RepositoryItemReader稳定读取所有符合初始过滤条件的页面,确保无数据遗漏。


解决方案

方案1:读取时加悲观锁,固定查询数据集

在查询阶段给数据加悲观锁,确保读取到的数据在处理完成前不会被修改,同时分页基于固定的数据集快照,避免状态变更干扰分页逻辑。

修改Repository查询方法

@Repository
public interface TransactionPagingAndSortingRepository extends PagingAndSortingRepository<Transaction, String> {
    @Lock(LockModeType.PESSIMISTIC_WRITE)
    Page<Transaction> findByCreatedAtBeforeAndStatusIn(ZonedDateTime currentDateTime, List<TransactionStatus> transactionStatuses, Pageable pageable);
}

说明:

  • PESSIMISTIC_WRITE锁会锁定查询到的行,直到当前事务提交,防止其他线程修改这些数据的状态。
  • 保持当前Reader中currenZonedDateTime的生成逻辑(Bean创建时一次性确定),避免后续页面查询因时间推移引入新数据,确保分页基于作业启动时的快照。

方案2:预取所有符合条件的ID,再分批读取实体

放弃直接分页查询实体,先一次性获取所有符合条件的事务ID,再基于固定的ID列表分批读取实体。这种方式的分页基准是静态ID集合,不受后续状态变更影响。

步骤1:新增Repository方法查询ID列表

@Repository
public interface TransactionPagingAndSortingRepository extends PagingAndSortingRepository<Transaction, String> {
    List<String> findIdsByCreatedAtBeforeAndStatusIn(ZonedDateTime currentDateTime, List<TransactionStatus> transactionStatuses);
    Transaction findById(String id);
}

步骤2:自定义ItemReader基于ID列表读取

@Bean
public ItemReader<Transaction> reader() {
    ZonedDateTime currentZonedDateTime = ZonedDateTime.now(ZoneId.of("UTC")).minusDays(1);
    List<TransactionStatus> transactionStatus = Collections.singletonList(TransactionStatus.NEW);
    
    // 作业启动时一次性获取所有符合条件的ID
    List<String> transactionIds = transactionPagingAndSortingRepository.findIdsByCreatedAtBeforeAndStatusIn(currentZonedDateTime, transactionStatus);
    
    // 基于ID列表分批读取实体
    return new ListItemReader<>(transactionIds) {
        @Override
        public Transaction read() throws Exception {
            String id = super.read();
            return id != null ? transactionPagingAndSortingRepository.findById(id) : null;
        }
    };
}

说明:

  • 若数据量极大(百万级以上),一次性查询ID列表会占用过多内存,可改为分页查询ID(比如按ID范围分批次获取)。
  • 需确保作业重启时能恢复ID列表的读取位置,可结合Spring Batch的状态保存机制实现。

方案3:改用主键范围分页,替代偏移量分页

将排序字段改为唯一且不变的主键id,通过主键范围查询代替传统的偏移量分页,避免状态变更导致的页面数据偏移。

修改Repository查询方法

@Repository
public interface TransactionPagingAndSortingRepository extends PagingAndSortingRepository<Transaction, String> {
    List<Transaction> findByCreatedAtBeforeAndStatusInAndIdGreaterThan(ZonedDateTime currentDateTime, List<TransactionStatus> transactionStatuses, String lastId, Pageable pageable);
}

自定义Reader实现范围分页逻辑

@Bean
public ItemReader<Transaction> reader() {
    ZonedDateTime currentZonedDateTime = ZonedDateTime.now(ZoneId.of("UTC")).minusDays(1);
    List<TransactionStatus> transactionStatus = Collections.singletonList(TransactionStatus.NEW);
    AtomicReference<String> lastId = new AtomicReference<>(null);
    int pageSize = 1000;
    List<Transaction> currentBatch = new ArrayList<>();

    return new ItemReader<>() {
        @Override
        public Transaction read() throws Exception {
            if (currentBatch.isEmpty()) {
                Pageable pageable = PageRequest.of(0, pageSize, Sort.by("id").ascending());
                if (lastId.get() == null) {
                    // 首次查询:获取ID最小的一批数据
                    currentBatch = transactionPagingAndSortingRepository.findByCreatedAtBeforeAndStatusIn(currentZonedDateTime, transactionStatus, pageable).getContent();
                } else {
                    // 后续查询:获取比上一批最后一个ID更大的数据
                    currentBatch = transactionPagingAndSortingRepository.findByCreatedAtBeforeAndStatusInAndIdGreaterThan(currentZonedDateTime, transactionStatus, lastId.get(), pageable);
                }
                if (!currentBatch.isEmpty()) {
                    lastId.set(currentBatch.get(currentBatch.size() - 1).getId());
                }
            }
            return currentBatch.isEmpty() ? null : currentBatch.remove(0);
        }
    };
}

说明:

  • 这种方式通过主键范围定位下一批数据,即使中间有数据状态变更,也不会跳过符合初始条件的数据。
  • 需确保主键是自增或有序的,保证范围查询的连续性。

额外注意事项

  • Writer逻辑补充:当前Writer为空实现,需补充数据库更新和Kafka消息发送逻辑,建议通过Spring事务管理保证更新与消息发送的原子性,避免数据不一致。
  • 作业状态兼容:若开启saveState(true),需确保状态保存逻辑与所选分页方式兼容(比如方案2中需保存ID列表的读取位置)。

内容的提问来源于stack exchange,提问作者Rohit Singh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 11:31:05