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

基于HTTP的Spring WebFlux背压逻辑疑问:我的预期是否正确?

问题描述

我有一个关于REST WebFlux端点在HTTP消费者速度较慢场景下预期行为的疑问。

我维护一个基于Java Spring Boot、使用WebFlux做数据流传输的API,有HTTP消费者调用该接口,消费者需要保持连接存活并持续读取数据。API每秒推送4000条事件,每条事件大小120字节。我们考虑到消费者可能会变慢,已经配置了背压机制来在这种情况下丢弃数据。

我们的预期行为:

  • 消费者仍能接收数据,只是速率变慢
  • 不会出现消费者完全停止接收数据的情况(假设消费者未阻塞)

实际现象:
存在正常传输时段(无背压触发的理想路径),也有消息丢失时段(符合预期)。但不符合预期的是:一旦出现消息丢失,此后客户端就不再接收任何数据,必须断开重连才能恢复数据流。我们原本预期:消息丢失可能发生,但客户端应仍能接收部分数据,且随着背压缓冲区排空,丢失情况会停止,无需重连即可恢复。

我们排查过客户端与服务器之间的网络问题,直接调用API(绕过所有反向代理)后仍出现相同现象,因此问题大概率出在REST API内部。

目前我们推测可能的原因:

  1. 代码实现有误或配置缺失/错误
  2. 相关库的最新版本已修复该问题
  3. 该行为是合理的,需要接受现状

恳请提供排查方向。

补充细节

API响应头

我认为客户端实现无关紧要。

< HTTP/1.1 200 OK
< Content-Type: application/json
< traceid: 0d5164d9c5771f1d
< Expires: Thu, 07 Nov 2024 13:33:55 GMT
< Cache-Control: max-age=0, no-cache, no-store
< Pragma: no-cache
< Date: Thu, 07 Nov 2024 13:33:55 GMT
< Transfer-Encoding:  chunked
< Connection: keep-alive
< Connection: Transfer-Encoding
< Strict-Transport-Security: max-age=31536000

API控制器

static final int BACKPRESSURE_BUFFER_MAX_SIZE = 8_000;

@GetMapping(value = "/streaming")
public Flux<?> streaming() {
   return dataStream.doFirst(() -> log.info("Streaming started"))
         .doOnNext(event -> counter.increment())
         .onBackpressureBuffer(BACKPRESSURE_BUFFER_MAX_SIZE, () -> droppedCounter.increment(), BufferOverflowStrategy.DROP_OLDEST)
         .doOnError(throwable -> log.error("Streaming error", throwable))
         .doFinally(signalType -> {
            log.info("Streaming finished. signalType={}", signalType);
         });
}

应用启动时的数据流配置

之后,在处理来自远程TCP服务器消息的消费者内部,使用tryEmitNext()方法填充sink。

var SCHEDULER = Schedulers.fromExecutorService(Executors.newFixedThreadPool(20));

var sink1 = Sinks.many().multicast().onBackpressureBuffer(10_000);
var sink2 = Sinks.many().multicast().onBackpressureBuffer(10_000);

var heartbeatStream =
      Flux.interval(Duration.ofSeconds(5))
            .map(sequence -> new HeartbeatResponse(Instant.now().getEpochSecond()));

var dataStream = Flux.merge(
            sink1.asFlux(),
            sink2.asFlux(),
            heartbeatStream)
      .publishOn(SCHEDULER)
      .share();

// dummy subscription done in @PostConstruct
dataStream.subscribe(message -> log.debug("Data received: {}", message));

依赖库版本

lib版本
spring boot2.7.18
spring webflux5.3.31
projectreactor core3.6.2
projectreactor netty1.1.15
netty4.1.101

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 12:28:12