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
相关产品推荐
相关产品推荐

