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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 20:21:34