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

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资源管理规则的错误写法共同导致:

  1. 手动干预背压请求逻辑,破坏响应式流语义
    代码中.doOnSubscribe(subscription -> subscription.request(64))是典型的错误写法。session.receive()返回的消息流本身由Spring WebFlux框架自动管理背压,下游publishOn、concatMap等操作符会根据自身实际消费能力,自动向上游请求对应数量的消息,无需手动调用request()。
    手动在订阅阶段一次性请求64条消息,会打乱整个链路的背压计数逻辑:当手动请求的64条消息消费完成后,下游操作符的自动请求逻辑被干扰,无法持续向上游拉取新消息,框架会判定接收端已无消费需求,直接触发onComplete信号关闭流。
  2. Netty引用计数对象内存泄漏
    代码中对WebSocketMessage调用了retain()增加引用计数,但消息处理完成后未调用release()归还内存。WebSocketMessage底层由Netty池化的直接内存ByteBuf支撑,未正常释放的ByteBuf会持续占用直接内存,当累积内存达到Netty内存池阈值时,底层会主动断开连接回收资源,该过程不会抛出业务层面的异常。
  3. 消费速度不匹配导致队列积压
    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 15:15:33