Reactor慢消费者仅处理最新值的实现疑问及问题排查
关于Reactor中onBackpressureLatest生效条件的疑问
在Reactor框架中,当存在快生产者+慢消费者的场景,且流中是快照类数据(比如GUI显示汇率的消费者、将汇率tick转为Flux的生产者)时,我们希望消费者只处理流中的最新值、丢弃旧值。此时Flux#onBackpressureLatest()算子看起来是合适的解决方案。
我搜索到的示例代码如下:
Flux.range(1, 30) .delayElements(Duration.ofMillis(500)) .onBackpressureLatest() .delayElements(Duration.ofMillis(3000)) .subscribe { println("got $it") }
这个示例在onBackpressureLatest()后手动添加了延迟,更接近Flux#sample(Duration)的效果,而非真实的慢消费者场景。
由于delayElements(Duration)内部封装了concatMap,我将代码改写为:
Flux.range(1, 30) .delayElements(Duration.ofMillis(500)) .onBackpressureLatest() .concatMap { Mono.just(it).subscribeOn(Schedulers.boundedElastic()) } .subscribe { item -> println("got $item") // 模拟慢消费者 Thread.sleep(3000) }
这个实现的思路和相关问题的答案一致,但我不理解为什么必须用concatMap(op)或者flatMap(op, 1, 1)才能让onBackpressureLatest()生效。
我尝试了以下简化版本,但都没达到预期效果,请问原因是什么?
// 无效尝试1 Flux.range(1, 30) .delayElements(Duration.ofMillis(500)) .onBackpressureLatest() .publishOn(Schedulers.boundedElastic()) .subscribe { item -> println("got $item") // 模拟慢消费者 Thread.sleep(3000) } // 无效尝试2 Flux.range(1, 30) .delayElements(Duration.ofMillis(500)) .onBackpressureLatest() .publishOn(Schedulers.boundedElastic()) .subscribe(object : BaseSubscriber<Int>() { override fun hookOnSubscribe(subscription: Subscription) { // 显式请求1条数据 subscription.request(1) } override fun hookOnNext(value: Int) { // 模拟慢消费者 Thread.sleep(3000) println("got $value") // 处理完后再请求1条 request(1) } })
内容的提问来源于stack exchange,提问作者Arie Xiao
相关产品推荐
相关产品推荐

