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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 21:20:29