Spring WebFlux WebSocket多聊天室消息定向推送配置问题
Spring响应式WebSocket按聊天室隔离消息实现方案
问题根因
当前实现使用全局单例的Sinks.Many<ChatMessage>存储所有消息,所有WebSocket连接无论属于哪个房间,都订阅同一个全局消息流,自然会收到所有聊天室的消息,没有做路径维度的消息隔离。
实现思路
- 从WebSocket握手请求的URI中提取当前连接对应的聊天室ID
- 为每个聊天室维护独立的消息Sink,不同房间的消息流完全隔离
- 移除手动
subscribe()的非规范写法,遵循响应式流生命周期由Spring容器管理 - 会话关闭时自动清理无连接的空房间Sink,避免内存泄漏
改造后代码
WebSocket处理器映射配置(原有配置可保留,无需修改)
@Bean public HandlerMapping webSocketHandlerMapping() { String path = "/chat/*"; Map<String, WebSocketHandler> map = Map.of(path, webSocketHandler); return new SimpleUrlHandlerMapping(map, -1); }
ChatSocketHandler改造实现
@Component public class ChatSocketHandler implements WebSocketHandler { private final ObjectMapper mapper = new ObjectMapper(); // 按房间ID存储独立的消息Sink,替换原全局单例Sink private final Map<String, Sinks.Many<ChatMessage>> roomSinks = new ConcurrentHashMap<>(); private final Sinks.EmitFailureHandler emitFailureHandler = (signalType, emitResult) -> emitResult.equals(Sinks.EmitResult.FAIL_NON_SERIALIZED); private final ChatService chatService; public ChatSocketHandler(ChatService chatService) { this.chatService = chatService; } @Override public Mono<Void> handle(WebSocketSession session) { // 从握手请求路径提取当前房间ID String requestPath = session.getHandshakeInfo().getUri().getPath(); String roomId = requestPath.substring(requestPath.lastIndexOf("/") + 1); // 初始化当前房间的消息Sink,已存在则直接复用 Sinks.Many<ChatMessage> currentRoomSink = roomSinks.computeIfAbsent(roomId, key -> Sinks.many().multicast().directBestEffort()); Flux<ChatMessage> currentRoomMessageFlux = currentRoomSink.asFlux(); // 入站消息处理:接收客户端消息、持久化、推送到当前房间的Sink Mono<Void> inbound = session.receive() .map(webSocketMessage -> { try { ChatMessage message = mapper.readValue(webSocketMessage.getPayloadAsText(), ChatMessage.class); // 强制绑定消息所属房间ID,防止跨房间发消息 message.setRoomId(roomId); return message; } catch (JsonProcessingException e) { e.printStackTrace(); return new ChatMessage(); } }) .flatMap(chatService::sendMessage) .doOnNext(message -> currentRoomSink.emitNext(message, emitFailureHandler)) .then(); // 出站消息处理:仅向当前客户端推送所属房间的消息 Mono<Void> outbound = session.send( currentRoomMessageFlux.map(message -> session.textMessage(toJson(message))) ); // 组合入站出站流,由框架管理流生命周期,无需手动subscribe return Mono.zip(inbound, outbound).then() .doFinally(signal -> { // 房间无任何连接时移除对应Sink,释放内存 if (currentRoomSink.currentSubscriberCount() == 0) { roomSinks.remove(roomId); } }); } private String toJson(ChatMessage object) { try { return mapper.writeValueAsString(object); } catch (JsonProcessingException e) { e.printStackTrace(); } return null; } }
注意事项
- 需确保
ChatMessage类包含roomId字段,持久化时同步存储房间ID,方便后续实现历史消息拉取功能 - 原实现中手动调用
subscribe()的写法不符合响应式编程规范,会导致流生命周期脱离容器管控,引发内存泄漏、连接无法正常释放问题,改造后通过Mono.zip组合流,由Spring WebSocket模块负责流的订阅和销毁 - 入站消息处理时强制覆盖消息的roomId,可避免客户端恶意构造请求跨房间发送消息
内容的提问来源于stack exchange,提问作者Jots
相关产品推荐
相关产品推荐

