WebFlux中onBackpressureError触发报错的条件及255阈值疑问
WebFlux onBackpressureError 触发报错与 requested=0 的关联解析
二者完全相关,你的观察是准确的,下面拆解具体逻辑:
核心流程与数字255的由来
你的代码是典型的快生产者+慢消费者场景:
- 上游
Flux.interval(Duration.ofMillis(1L))每1毫秒生成一条数据,生产速度极快; - 下游subscribe里用
Thread.sleep(5L)模拟耗时处理,消费速度远慢于生产。
中间的publishOn(Schedulers.parallel())有个关键细节:它默认的缓冲区(prefetch)大小是256。下游订阅时,publishOn会先向上游请求256条数据,这些数据会先存入缓冲区,再慢慢交给下游处理。
但因为消费速度慢,缓冲区很快被填满。当上游生产到第256条数据(对应Flux.interval的索引值255,从0开始计数)时,缓冲区已经没有剩余空间,publishOn会停止向上游发起新的请求,此时DropSubscriber(onBackpressureError的内部订阅者)的requested计数器就会减到0。
onBackpressureError 的触发逻辑
onBackpressureError的作用就是:当下游不再请求数据(即requested=0)时,上游如果继续发送数据,直接抛出背压异常。
反编译看到的DropSubscriber.requested变为0,就是触发报错的直接条件:
- 当
requested > 0时,DropSubscriber会正常接收数据并转发,同时将requested减1; - 当
requested == 0时,上游再调用onNext,DropSubscriber就会立即抛出IllegalStateException,终止流。
验证你的现象
你看到前255条数据正常接收,是因为这对应了索引0到254的255条数据,此时缓冲区还剩1个空位;当第256条数据(索引255)过来时,缓冲区填满,requested耗尽变为0,触发onBackpressureError的报错逻辑。
如果想调整这个阈值,可以修改publishOn的prefetch参数,比如publishOn(Schedulers.parallel(), 1024),此时报错会推迟到第1025条数据。
内容的提问来源于stack exchange,提问作者Joonseo Lee
相关产品推荐
相关产品推荐

