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

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

批量请求相比逐个请求确实有明显优势:

  1. 减少交互开销:降低了request()的调用次数,减少了生产者和消费者之间的信号传递次数。对于跨线程或跨进程的异步数据流,能减少上下文切换和调度的额外开销。
  2. 提升吞吐量:生产者可以利用批量处理的优势,比如数据库批量查询、IO批量读取,避免重复初始化和销毁资源,提高数据生成或获取的效率。
  3. 优化端到端延迟:如果生产者需要等待外部资源(比如网络响应),批量请求可以让生产者提前发起批量获取,避免每次处理完一个值才发起请求的等待时间,整体延迟会更低。

不过批量请求也有适用边界:如果消费者处理单个数据的时间波动较大,批量请求可能导致内存中堆积过多未处理数据,增加内存压力。这种情况下需要根据实际处理能力调整批量大小,或者结合动态背压策略。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 06:31:15