Spring WebFlux WebSocket接收消息过快时无异常断开连接问题
问题代码
public Mono<Void> handle(@Nonnull WebSocketSession session) { final WebSocketContext webSocketContext = new WebSocketContext(session); Mono<Void> output = session.send(Flux.create(webSocketContext::setSink)); Mono<Void> input = session.receive() .timeout(Duration.ofSeconds(adapterProperties.getSessionTimeout())) .doOnSubscribe(subscription -> subscription.request(64)) .doOnNext(WebSocketMessage::retain) .publishOn(Schedulers.boundedElastic()) .concatMap(msg -> { // ....blocking operation return Flux.empty(); }).then(); return Mono.zip(input, output).then(); }
问题现象
- WebSocket客户端高速向服务端发送消息时,服务端累计接收约2000条数据后连接自动断开,全程无异常信息抛出
- 调慢客户端消息发送速率后,连接可长期稳定保持正常运行
已采集运行日志
2022-07-14 17:19:40.295 adapter-iat [boundedElastic-5] INFO reactor.Flux.PublishOn.5 - | onNext(WebSocket TEXT message (13765 bytes)) 2022-07-14 17:19:40.296 adapter-iat [boundedElastic-5] INFO reactor.Flux.PublishOn.5 - | onNext(WebSocket TEXT message (13765 bytes)) 2022-07-14 17:19:40.300 adapter-iat [boundedElastic-5] INFO reactor.Flux.PublishOn.5 - | onComplete()
上述
onComplete信号无明确业务触发逻辑,触发原因未知
排查思路与解决方案
核心根因
该问题由三个违反响应式编程和Netty资源管理规则的错误写法共同导致:
- 手动干预背压请求逻辑,破坏响应式流语义
代码中.doOnSubscribe(subscription -> subscription.request(64))是典型的错误写法。session.receive()返回的消息流本身由Spring WebFlux框架自动管理背压,下游publishOn、concatMap等操作符会根据自身实际消费能力,自动向上游请求对应数量的消息,无需手动调用request()。
手动在订阅阶段一次性请求64条消息,会打乱整个链路的背压计数逻辑:当手动请求的64条消息消费完成后,下游操作符的自动请求逻辑被干扰,无法持续向上游拉取新消息,框架会判定接收端已无消费需求,直接触发onComplete信号关闭流。 - Netty引用计数对象内存泄漏
代码中对WebSocketMessage调用了retain()增加引用计数,但消息处理完成后未调用release()归还内存。WebSocketMessage底层由Netty池化的直接内存ByteBuf支撑,未正常释放的ByteBuf会持续占用直接内存,当累积内存达到Netty内存池阈值时,底层会主动断开连接回收资源,该过程不会抛出业务层面的异常。 - 消费速度不匹配导致队列积压
concatMap中执行阻塞操作,虽然切换到了boundedElastic调度器,但concatMap本身是按顺序串行处理消息,阻塞操作的吞吐量远低于客户端高速发送的速度时,消息会在各个操作符的内部队列、Netty TCP接收缓冲区持续积压,当积压量超过缓冲区阈值时,TCP层会直接断开连接。
修复步骤
- 直接删除
.doOnSubscribe(subscription -> subscription.request(64))代码,完全交由响应式链路的操作符自动管理背压请求,非自定义操作符场景下禁止手动调用Subscription的request()方法。 - 修复内存泄漏:在
concatMap的阻塞处理逻辑执行完成后,必须调用msg.release()释放消息占用的内存;同时新增doOnDiscard(WebSocketMessage.class, WebSocketMessage::release)逻辑,统一处理队列溢出、流取消场景下被丢弃的消息,避免ByteBuf泄漏。 - 优化背压匹配逻辑:如果阻塞操作的吞吐量确实无法匹配客户端发送速率,给
publishOn、concatMap配置合理的预取参数,同时根据业务需求添加onBackpressureDrop/onBackpressureLatest策略,避免无限制缓冲消息导致内存溢出、连接断连。 - 补充全链路信号监听:在
input和output流上新增doOnError、doOnCancel、doOnTerminate的日志打印,记录错误栈、取消信号触发点,避免异常被操作符吞掉导致无问题现场。 - 给发送端的
Flux.create显式配置背压策略,不要使用默认的无界缓冲策略,避免发送端消息积压导致内存溢出。
内容的提问来源于stack exchange,提问作者nobody
相关产品推荐
相关产品推荐

