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
失败原因与正确实现
失败原因
初始写法存在两个核心错误:
- 操作符顺序错误:
publishOn内部自带默认容量256的预取队列,会提前从上游拉取数据缓存。onBackpressureLatest()放在publishOn下游,无法拦截已经被publishOn预取进内部队列的旧值,背压策略完全不生效。 limitRate(1)不符合预期行为:Reactor的limitRate默认自带预取优化,即使传入参数1,内部也会提前批量拉取数据存入内部队列,所有中间产生的值都会被缓存后挨个交给消费者,无法实现丢弃旧值的效果。
正确实现方案
调整操作符顺序与配置即可:
- 将
onBackpressureLatest()放在最靠近上游生产者的位置,且位于publishOn之前,让背压丢弃策略直接作用于生产端,仅保留最新值,不缓存所有中间值。 - 移除
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
相关产品推荐
相关产品推荐

