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

Spring WebFlux中WebSocketSession.send()无法发送消息问题排查

问题核心原因

你代码的问题出在EventService中Flux.create的实现逻辑,无限死循环阻塞了WebSocket的IO处理线程,导致即使生成了事件也无法正常发送到客户端。

具体问题拆解

  • 阻塞线程占用:Flux.create默认会在订阅触发的线程(也就是WebSocket请求的处理线程)中执行你传入的生成逻辑,你的while(true)没有任何休眠/让出CPU的逻辑,一直占用当前线程空转,Reactor的非阻塞IO模型下,线程被占死后后续的消息发送、网络IO操作根本得不到执行,而循环内的println是在死循环内直接执行的,所以控制台可以正常打印。
  • 线程安全问题:你用来暂存事件的ArrayList是线程不安全的,当其他线程调用push方法写入事件、同时Flux.create的循环在读取/删除元素时,会出现并发异常、数据丢失问题。
  • Flux实现逻辑不合理:手动通过死循环拉取事件的实现完全不符合Reactor的设计规范,也没有处理背压场景,很容易出现OOM或者数据丢失。

修复方案

推荐用Sinks来实现多播事件流,无需自己写循环处理事件生成逻辑,代码如下:

修正后的EventService

@Service
public class EventService {
    // 构造多播的Sink,支持背压缓存
    private final Sinks.Many<EventDto> eventSink = Sinks.many().multicast().onBackpressureBuffer();
    private final Flux<EventDto> eventFlux = eventSink.asFlux().publish().autoConnect();

    public void push(EventDto event) {
        // 发送事件,推送结果的异常处理可按自己需求调整
        eventSink.tryEmitNext(event);
    }

    public Flux<EventDto> events() {
        return eventFlux;
    }
}

优化后的EventWebsocketHandler

public class EventWebsocketHandler implements WebSocketHandler {
    // 复用ObjectMapper,不要每次handle都新建
    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
    private final EventService eventService;

    public EventWebsocketHandler(EventService eventService) {
        this.eventService = eventService;
    }

    @Override
    public Mono<Void> handle(WebSocketSession session) {
        Flux<WebSocketMessage> messages = eventService.events()
                .handle((event, sink) -> {
                    try {
                        System.out.println(event);
                        sink.next(OBJECT_MAPPER.writeValueAsString(event));
                    } catch (JsonProcessingException e) {
                        sink.error(e);
                    }
                })
                .map(session::textMessage);
        return session.send(messages);
    }
}

修正后不需要死循环占用线程,Sinks会自动处理事件的推送,线程不会被阻塞,WebSocket的消息发送逻辑就可以正常执行,客户端就能收到消息了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 18:57:03