如何让Flux始终提前请求一个元素,在订阅者处理数据时预取下一项
根因分析
你当前使用的Flux.generate是同步按需生成的算子,默认每次仅响应下游的1个元素请求,只有下游处理完当前元素、发起新的拉取请求后,才会触发下一次generate回调执行hasNext()和next()逻辑,本身不支持主动预生成元素,这也是你单独加publish(2)没有达到预期的核心原因。
解决方案
方案1:改用Flux.fromIterable适配预取逻辑
Flux.fromIterable原生支持背压预取能力,你只需要将自定义Iterator包装为Iterable,再配合limitRate指定预取数即可,实现成本最低:
@Override public Flux<SearchResponse> createDeepQueryFlux(@NonNull PITSearchInput input) { validator.validate(input); // 自定义Iterator包装为Iterable Iterable<SearchResponse> pitIterable = () -> new PointInTimeIterator(forwardingRestHighLevelClient, input); return Flux.fromIterable(pitIterable) // 预取数设置为2,实现提前拉取下一个元素的效果 .limitRate(2) .subscribeOn(Schedulers.boundedElastic()) // 处理Iterator资源关闭逻辑 .doFinally(signalType -> { if (pitIterable.iterator() instanceof PointInTimeIterator pitIterator) { pitIterator.shutdown(); } }); }
方案2:保留Flux.generate,搭配publishOn预取参数
如果你需要保留generate的状态管理、资源关闭逻辑,可以通过publishOn的预取参数实现主动拉取:publishOn会提前向上游请求指定数量的元素缓存,下游处理逻辑运行在publishOn指定的线程池时,生成逻辑就可以在boundedElastic线程中提前生成下一个元素,和下游处理并行执行:
@Override public Flux<SearchResponse> createDeepQueryFlux(@NonNull PITSearchInput input) { validator.validate(input); return Flux.<SearchResponse, PointInTimeIterator>generate( () -> new PointInTimeIterator(forwardingRestHighLevelClient, input), (deepQueryIterator, sink) -> { log.info("Generate called"); if (deepQueryIterator.hasNext()) { sink.next(deepQueryIterator.next()); } else { sink.complete(); } return deepQueryIterator; }, PointInTimeIterator::shutdown ) .subscribeOn(Schedulers.boundedElastic()) // 第二个参数为预取数,设为2即提前缓存2个元素 .publishOn(Schedulers.parallel(), 2); }
效果验证
你可以在PointInTimeIterator的next()方法中增加日志,同时在订阅者的处理逻辑中增加sleep模拟业务耗时,就能看到next()的调用时间早于订阅者处理完上一个元素的时间,符合预取需求。
内容的提问来源于stack exchange,提问作者John B
相关产品推荐
相关产品推荐

