如何用Reactive Streams实现带超时与提前退出的轮询器?
没问题,我来帮你解决这两个Reactive Streams轮询器的实现问题,结合Java生态里常用的Reactor框架(Spring体系标配)给你具体的代码方案和思路:
1. 实现带超时的Reactive Streams轮询器
带超时的轮询器核心需求是:按固定间隔重复查询数据,直到超时时间到达,然后返回超时前的最新结果。我们可以利用Reactor的Flux.interval生成轮询触发信号,结合缓存、超时终止操作符来实现:
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.time.Duration; public class TimeoutPoller { // 模拟实际的轮询数据获取逻辑(比如调用外部API、查询数据库) private Mono<Integer> fetchPollData() { // 这里替换为你的真实业务逻辑 return Mono.just((int) (Math.random() * 100)); } /** * 带超时的轮询器方法 * @param pollInterval 轮询间隔 * @param timeout 超时时间 * @return 超时前的最新轮询结果 */ public Mono<Integer> startPollWithTimeout(Duration pollInterval, Duration timeout) { return Flux.interval(pollInterval) .flatMap(ignored -> fetchPollData()) // 每次间隔触发时执行数据查询 .cache() // 缓存最新的查询结果,确保超时后能拿到最后一次有效数据 .takeUntilOther(Mono.delay(timeout)) // 超时时间到后立即终止轮询流 .last() // 取出流终止前的最后一个数据(即最新轮询结果) .timeout(timeout); // 双重保障:如果首次轮询就卡住,直接触发超时 } }
关键逻辑说明:
Flux.interval(pollInterval):按指定间隔生成递增的Long信号,作为轮询的触发开关cache():缓存流中最新的元素,避免超时后丢失最后一次的轮询结果takeUntilOther(Mono.delay(timeout)):当延迟Mono触发(超时时间到),立刻终止轮询流last():专门获取流终止前的最后一个元素,完美契合“取最新事件”的需求- 额外的
timeout(timeout):防止轮询逻辑长时间无响应(比如首次查询就阻塞),直接触发超时错误
2. 建模带超时+提前退出条件的轮询器(结合你的示例)
针对你提到的场景:有一个每秒生成9,8,7,...,0的序列,需要在5秒后取最新事件,或者在0提前出现时直接取0并退出。我们可以结合takeUntil(满足条件提前终止)和takeUntilOther(超时终止)来实现双条件终止逻辑:
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.time.Duration; import java.util.concurrent.atomic.AtomicInteger; public class EarlyExitPoller { // 模拟你描述的每秒生成递减序列的Flux private Flux<Integer> generateDecreasingSequence() { AtomicInteger counter = new AtomicInteger(9); return Flux.interval(Duration.ofSeconds(1)) .map(ignored -> counter.getAndDecrement()) .takeWhile(num -> num >= 0); // 生成到0后自动停止序列 } /** * 带超时和提前退出条件的消费者方法 * @return 要么是提前出现的0,要么是5秒时的最新序列值 */ public Mono<Integer> consumeSequence() { Duration timeout = Duration.ofSeconds(5); return generateDecreasingSequence() .takeUntil(num -> num == 0) // 一旦出现0,立即终止流 .takeUntilOther(Mono.delay(timeout)) // 5秒超时后终止流 .last() // 取出终止时的最后一个元素 .defaultIfEmpty(-1); // 兜底:处理流未产生任何元素的极端情况 } }
关键逻辑说明:
takeUntil(num -> num == 0):这是提前退出的核心,只要流中出现0,立刻终止整个流,优先响应业务条件takeUntilOther(Mono.delay(timeout)):和提前终止条件并行生效,5秒时间一到,不管序列有没有到0,都终止流last():自动判断哪个条件先触发——如果0在5秒前出现,就返回0;如果5秒到了序列还没到0(比如此时序列到了5),就返回5defaultIfEmpty(-1):处理极端场景(比如刚启动就超时),避免返回空值引发异常
内容的提问来源于stack exchange,提问作者user2561460
相关产品推荐
相关产品推荐

