基于HTTP的Spring WebFlux背压逻辑疑问:我的预期是否正确?
问题描述
我有一个关于REST WebFlux端点在HTTP消费者速度较慢场景下预期行为的疑问。
我维护一个基于Java Spring Boot、使用WebFlux做数据流传输的API,有HTTP消费者调用该接口,消费者需要保持连接存活并持续读取数据。API每秒推送4000条事件,每条事件大小120字节。我们考虑到消费者可能会变慢,已经配置了背压机制来在这种情况下丢弃数据。
我们的预期行为:
- 消费者仍能接收数据,只是速率变慢
- 不会出现消费者完全停止接收数据的情况(假设消费者未阻塞)
实际现象:
存在正常传输时段(无背压触发的理想路径),也有消息丢失时段(符合预期)。但不符合预期的是:一旦出现消息丢失,此后客户端就不再接收任何数据,必须断开重连才能恢复数据流。我们原本预期:消息丢失可能发生,但客户端应仍能接收部分数据,且随着背压缓冲区排空,丢失情况会停止,无需重连即可恢复。
我们排查过客户端与服务器之间的网络问题,直接调用API(绕过所有反向代理)后仍出现相同现象,因此问题大概率出在REST API内部。
目前我们推测可能的原因:
- 代码实现有误或配置缺失/错误
- 相关库的最新版本已修复该问题
- 该行为是合理的,需要接受现状
恳请提供排查方向。
补充细节
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 boot | 2.7.18 |
| spring webflux | 5.3.31 |
| projectreactor core | 3.6.2 |
| projectreactor netty | 1.1.15 |
| netty | 4.1.101 |
内容的提问来源于stack exchange,提问作者Michal
相关产品推荐
相关产品推荐

