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

Spring Reactor背压场景下仅保留最新值覆盖旧值的实现方案

问题场景

使用Project Spring Reactor实现典型响应式业务逻辑:生产者数据生成速度快于消费者处理速度,当新值可用时,消费者无需处理旧值(典型案例:过时的股票价格无业务价值,无需消费)。
本次测试的固定参数:

  • 生产者每100ms生成一个新值
  • 消费者处理单个值耗时500ms
  • 目标效果:消费者两次处理的间隔期内产生的多个新值,仅消费最新值,忽略所有过时的中间值
    初始尝试方案:通过配置limitRate(1)实现每次仅向生产者请求1个值,搭配onBackpressureLatest()达到忽略中间值的效果,但组合使用未达到预期。
初始测试代码
@Test
void fluxTest(){
    
    Flux<Integer> flux = Flux.generate(AtomicInteger::new, (ai, sink) -> {

        int i = ai.incrementAndGet();

        if (i > 10) {
            sink.complete();
        } else {
            System.out.println(Thread.currentThread()+": generate & emit value "+i);
            sink.next(i);
        }
        sleep(100);
        return ai;
    });

    flux
            .publishOn(Schedulers.parallel())
            .onBackpressureLatest()
            .limitRate(1)
            .subscribe(i -> {
                System.out.println(Thread.currentThread()+": Receive: " + i); // do something with generated and processed item
                sleep(500);
            });

    sleep(10000);
}

void sleep(int ms){
    try {
        Thread.sleep(ms);
    } catch (InterruptedException e) {
        e.printStackTrace();
    }
}
实际运行结果
14:26:40.019 [main] DEBUG reactor.util.Loggers - Using Slf4j logging framework
Thread[main,5,main]: generate & emit value 1
Thread[parallel-1,5,main]: Receive: 1
Thread[main,5,main]: generate & emit value 2
Thread[main,5,main]: generate & emit value 3
Thread[main,5,main]: generate & emit value 4
Thread[main,5,main]: generate & emit value 5
Thread[parallel-1,5,main]: Receive: 2
Thread[main,5,main]: generate & emit value 6
Thread[main,5,main]: generate & emit value 7
Thread[main,5,main]: generate & emit value 8
Thread[main,5,main]: generate & emit value 9
Thread[main,5,main]: generate & emit value 10
Thread[parallel-1,5,main]: Receive: 3
Thread[parallel-1,5,main]: Receive: 4
Thread[parallel-1,5,main]: Receive: 5
Thread[parallel-1,5,main]: Receive: 6
Thread[parallel-1,5,main]: Receive: 7
Thread[parallel-1,5,main]: Receive: 8
Thread[parallel-1,5,main]: Receive: 9
Thread[parallel-1,5,main]: Receive: 10

Process finished with exit code 0
预期运行结果
14:26:40.019 [main] DEBUG reactor.util.Loggers - Using Slf4j logging framework
Thread[main,5,main]: generate & emit value 1
Thread[parallel-1,5,main]: Receive: 1
Thread[main,5,main]: generate & emit value 2
Thread[main,5,main]: generate & emit value 3
Thread[main,5,main]: generate & emit value 4
Thread[main,5,main]: generate & emit value 5
Thread[parallel-1,5,main]: Receive: 5
Thread[main,5,main]: generate & emit value 6
Thread[main,5,main]: generate & emit value 7
Thread[main,5,main]: generate & emit value 8
Thread[main,5,main]: generate & emit value 9
Thread[main,5,main]: generate & emit value 10
Thread[parallel-1,5,main]: Receive: 10

Process finished with exit code 0
失败原因与正确实现

失败原因

初始写法存在两个核心错误:

  1. 操作符顺序错误:publishOn内部自带默认容量256的预取队列,会提前从上游拉取数据缓存。onBackpressureLatest()放在publishOn下游,无法拦截已经被publishOn预取进内部队列的旧值,背压策略完全不生效。
  2. limitRate(1)不符合预期行为:Reactor的limitRate默认自带预取优化,即使传入参数1,内部也会提前批量拉取数据存入内部队列,所有中间产生的值都会被缓存后挨个交给消费者,无法实现丢弃旧值的效果。

正确实现方案

调整操作符顺序与配置即可:

  1. 将onBackpressureLatest()放在最靠近上游生产者的位置,且位于publishOn之前,让背压丢弃策略直接作用于生产端,仅保留最新值,不缓存所有中间值。
  2. 移除limitRate(1),直接给publishOn配置预取参数:队列大小设为1,关闭提前预取优化,让消费者每处理完1个值,才向上游请求1个新值。

修正后的核心链式调用代码:

flux
        .onBackpressureLatest() // 背压策略紧贴上游,直接对生产端生效
        .publishOn(Schedulers.parallel(), false, 1) // 预取数设为1,关闭预取优化
        .subscribe(i -> {
            System.out.println(Thread.currentThread()+": Receive: " + i);
            sleep(500);
        });

运行后即可得到预期结果:消费者每500ms处理一次,每次都能拿到当前最新生成的值,所有中间旧值会被自动丢弃。


内容的提问来源于stack exchange,提问作者SvenR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 13:51:10