Java运行时将分页查询接口封装为惰性加载Stream的实现问题
分页场景惰性Stream实现方案
问题1答复
可以实现,完全不需要提前知道元素总数。Java Stream本身基于惰性求值设计,我们可以通过自定义Spliterator的方式实现分页数据的按需加载,调用streamOf方法时不会触发任何接口请求,只有在真正遍历Stream元素时才会逐页拉取数据。
问题2可行实现方案
方案1:无需提前获取总数的通用实现(推荐)
该实现适配所有分页接口,无需依赖Page返回总条数,只要当offset超过总元素数时返回空列表即可正常运行:
import java.util.Collections; import java.util.Iterator; import java.util.Spliterators; import java.util.function.BiFunction; import java.util.function.Consumer; import java.util.stream.Stream; import java.util.stream.StreamSupport; public class PageStreamUtils { // 可自定义默认页大小 private static final int DEFAULT_PAGE_SIZE = 20; public static <T> Stream<T> streamOf(BiFunction<Integer, Integer, Page<T>> pageFunction) { return streamOf(pageFunction, DEFAULT_PAGE_SIZE); } public static <T> Stream<T> streamOf(BiFunction<Integer, Integer, Page<T>> pageFunction, int pageSize) { if (pageSize <= 0) { throw new IllegalArgumentException("pageSize must be positive"); } return StreamSupport.stream( new Spliterators.AbstractSpliterator<T>(Long.MAX_VALUE, Spliterator.ORDERED) { private int currentOffset = 0; private Iterator<T> currentPageIter = Collections.emptyIterator(); @Override public boolean tryAdvance(Consumer<? super T> action) { // 当前页元素遍历完后加载下一页 while (!currentPageIter.hasNext()) { Page<T> page = pageFunction.apply(currentOffset, pageSize); if (page.getContent().isEmpty()) { // 无更多数据,终止遍历 return false; } currentPageIter = page.getContent().iterator(); currentOffset += pageSize; } action.accept(currentPageIter.next()); return true; } }, // 串行流,分页接口通常不建议并行调用避免触发限流 false ); } }
使用方式和预期完全一致:
Stream<User> userStream = PageStreamUtils.streamOf(service::users); // 此时还没有发起任何接口请求,遍历的时候才会逐页拉取 userStream.forEach(System.out::println);
方案2:提前获取总条数的简化实现
如果你的Page对象一定会返回总条数,且需要提前拿到总元素数做业务校验,可以用如下实现,仅在第一次遍历的时候发起一次总条数查询:
public static <T> Stream<T> streamOfWithTotal(BiFunction<Integer, Integer, Page<T>> pageFunction, int pageSize) { if (pageSize <= 0) { throw new IllegalArgumentException("pageSize must be positive"); } return StreamSupport.stream( new Spliterators.AbstractSpliterator<T>(Long.MAX_VALUE, Spliterator.SIZED | Spliterator.ORDERED) { private int currentOffset = 0; private Iterator<T> currentPageIter = Collections.emptyIterator(); private Long total = null; @Override public boolean tryAdvance(Consumer<? super T> action) { while (!currentPageIter.hasNext()) { if (total != null && currentOffset >= total) { return false; } Page<T> page = pageFunction.apply(currentOffset, pageSize); if (total == null) { total = page.getTotalElements(); if (total == 0) { return false; } } if (page.getContent().isEmpty()) { return false; } currentPageIter = page.getContent().iterator(); currentOffset += pageSize; } action.accept(currentPageIter.next()); return true; } @Override public long estimateSize() { if (total == null) { // 还未查询总条数时返回预估大小 return Long.MAX_VALUE; } return total; } }, false ); }
注意事项
- 生成的Stream默认是串行流,不建议并行处理,避免短时间内大量请求打垮下游分页接口
- 如果分页接口有调用频率限制,可以在加载下一页的逻辑中加入限流逻辑
- Stream使用完会自动释放相关资源,不需要额外关闭,如果你有特殊资源释放需求,可以通过
Stream.onClose()方法添加回调
内容的提问来源于stack exchange,提问作者Alexandr
相关产品推荐
相关产品推荐

