Reactor背压:hookOnNext中多request的作用及批量请求优势问询
Reactor背压机制测试与疑问解答
我正在测试Reactor的背压机制,以下是常见的基础实现结构:
package com.example.backpressure; import org.junit.jupiter.api.Test; import org.reactivestreams.Subscription; import org.springframework.boot.test.context.SpringBootTest; import reactor.core.publisher.BaseSubscriber; import reactor.core.publisher.Flux; @SpringBootTest class BackPressureApplicationTests { @Test public void backPressureApplicationTest1() { Flux<Integer> request = Flux.range(1, 20).log(); request.subscribe(new BackPressureSubscriber<>()); } } class BackPressureSubscriber<T> extends BaseSubscriber<T> { public void hookOnSubscribe(Subscription subscription) { request(1); } public void hookOnNext(T value) { System.out.println("Value is: " + value); request(1); } }
该实现通过处理完当前值后请求1个新值(request(1))施加背压,运行输出符合预期:
2022-10-03 20:11:51.908 INFO 2380 --- [ main] reactor.Flux.Range.1 : | onSubscribe([Synchronous Fuseable] FluxRange.RangeSubscription) 2022-10-03 20:11:51.909 INFO 2380 --- [ main] reactor.Flux.Range.1 : | request(1) 2022-10-03 20:11:51.910 INFO 2380 --- [ main] reactor.Flux.Range.1 : | onNext(1) Value is: 1 2022-10-03 20:11:51.911 INFO 2380 --- [ main] reactor.Flux.Range.1 : | request(1) 2022-10-03 20:11:51.911 INFO 2380 --- [ main] reactor.Flux.Range.1 : | onNext(2) Value is: 2 2022-10-03 20:11:51.911 INFO 2380 --- [ main] reactor.Flux.Range.1 : | request(1) 2022-10-03 20:11:51.911 INFO 2380 --- [ main] reactor.Flux.Range.1 : | onNext(3) Value is: 3 2022-10-03 20:11:51.911 INFO 2380 --- [ main] reactor.Flux.Range.1 : | request(1) 2022-10-03 20:11:51.911 INFO 2380 --- [ main] reactor.Flux.Range.1 : | onNext(4) Value is: 4
若修改hookOnNext方法,每次处理完值后调用request(3):
public void hookOnNext(T value) { System.out.println("Value is: " + value); request(3); }
输出仅request数值变化,由此提出两个技术疑问,以下是解答:
疑问1:每次接收值后请求3个新值,生产者的请求配额会累加吗?在hookOnNext中请求超过1个值是否有实际意义?
- 请求配额会累加:Reactor的
Subscription内部维护一个未完成的请求配额计数器,每次调用request(n)都会把n加到当前配额里。比如第一次请求1,处理完第一个值后再请求3,此时配额变为3,生产者会一次性推送3个值,直到配额耗尽。 - 请求超过1个值有实际意义:这种方式能减少生产者和消费者之间的请求信号交互次数。如果生产者需要从外部系统获取数据(比如数据库查询、网络请求),批量请求可以让生产者一次性获取多个数据,降低IO或资源调度的开销,提升整体吞吐量。但要注意,如果消费者处理速度跟不上,可能会导致内存中堆积过多未处理数据,削弱背压的节流效果。
疑问2:每处理3个值后请求3个新值的批量方式,相比逐个请求是否存在优势,比如降低延迟?
先看对应的批量请求实现:
class BackPressureSubscriber<T> extends BaseSubscriber<T> { int cnt = 0; int numEachTime=3; public void hookOnSubscribe(Subscription subscription) { request(3); } public void hookOnNext(T value) { System.out.println("Value is: " + value); cnt++; if(cnt%numEachTime == 0) { request(3); } } }
运行输出显示每处理3个值后批量推送新值:
2022-10-03 20:21:59.681 INFO 11836 --- [ main] reactor.Flux.Range.1 : | onSubscribe([Synchronous Fuseable] FluxRange.RangeSubscription) 2022-10-03 20:21:59.683 INFO 11836 --- [ main] reactor.Flux.Range.1 : | request(3) 2022-10-03 20:21:59.683 INFO 11836 --- [ main] reactor.Flux.Range.1 : | onNext(1) Value is: 1 2022-10-03 20:21:59.683 INFO 11836 --- [ main] reactor.Flux.Range.1 : | onNext(2) Value is: 2 2022-10-03 20:21:59.683 INFO 11836 --- [ main] reactor.Flux.Range.1 : | onNext(3) Value is: 3 2022-10-03 20:21:59.683 INFO 11836 --- [ main] reactor.Flux.Range.1 : | request(3) 2022-10-03 20:21:59.683 INFO 11836 --- [ main] reactor.Flux.Range.1 : | onNext(4) Value is: 4 2022-10-03 20:21:59.683 INFO 11836 --- [ main] reactor.Flux.Range.1 : | onNext(5) Value is: 5 2022-10-03 20:21:59.683 INFO 11836 --- [ main] reactor.Flux.Range.1 : | onNext(6) Value is: 6
批量请求相比逐个请求确实有明显优势:
- 减少交互开销:降低了
request()的调用次数,减少了生产者和消费者之间的信号传递次数。对于跨线程或跨进程的异步数据流,能减少上下文切换和调度的额外开销。 - 提升吞吐量:生产者可以利用批量处理的优势,比如数据库批量查询、IO批量读取,避免重复初始化和销毁资源,提高数据生成或获取的效率。
- 优化端到端延迟:如果生产者需要等待外部资源(比如网络响应),批量请求可以让生产者提前发起批量获取,避免每次处理完一个值才发起请求的等待时间,整体延迟会更低。
不过批量请求也有适用边界:如果消费者处理单个数据的时间波动较大,批量请求可能导致内存中堆积过多未处理数据,增加内存压力。这种情况下需要根据实际处理能力调整批量大小,或者结合动态背压策略。
内容的提问来源于stack exchange,提问作者Ziggy000
相关产品推荐
相关产品推荐

