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

WebSocket会话订阅指定会话消息的MongoDB轮询实现优化咨询

问题描述

我有一个WebSocket端点,允许客户端订阅特定会话(conversation)相关的消息。由于服务采用多实例部署,消息存储在MongoDB集合中,因此每个WebSocket会话需要启动一个轮询器,仅查询与其关联的消息。

同时,可能存在多个WebSocket会话订阅同一个会话(conversation)的情况。

WebSocket客户端需在握手时提供conversationId查询参数(也可以作为URL路径的一部分):

ws://my.api.com/wsflow/websocket?conversationId=12345

我的实现逻辑:

  • 创建ServerWebSocketContainer,设置HandshakeInterceptor获取conversationId并保存为会话属性
  • 通过WebSocketHandlerDecorator重写afterConnectionEstablished()方法,为每个WebSocket会话创建一个集成流,轮询对应conversationId的消息
  • 为避免重复推送消息给同一会话的多个订阅者,在MongoDB文档中添加sessions数组:轮询器仅查询当前WS会话ID不在该数组中的消息,同时通过push操作将当前会话ID追加到数组中
  • 查询到的消息会填充sessionId头,发送到wsMessages通道,由WebSocketOutboundMessageHandler推送给客户端

目前代码在本地可以正常运行,但存在两个疑问:

  1. 是否有更简洁自然的实现方式?
  2. useFlowIdAsPrefix()是否必要?它会对afterConnectionClosed()中执行的flowContext.remove()产生影响吗?
实现代码
@Bean
public ServerWebSocketContainer serverContainer(MongoTemplate mongoTemplate, IntegrationFlowContext flowContext) {
    ServerWebSocketContainer container = new ServerWebSocketContainer("/wsflow").withSockJs(
        new SockJsServiceOptions().setHeartbeatTime(5000L).setTaskScheduler(new TaskSchedulerBuilder().build()));
    container.setInterceptors(new HandshakeInterceptor() {
        @Override
        public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response,
                WebSocketHandler wsHandler, Exception exception) {}
        @Override
        public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response,
                WebSocketHandler wsHandler, Map<String, Object> attributes) throws Exception {
            String conversationId = ((ServletServerHttpRequest)request).getServletRequest().getParameter("conversationId");
            attributes.put("conversationId", conversationId);       
            return true;
        }
    });
    container.setDecoratorFactories(handler -> new WebSocketHandlerDecorator(handler) {
        @Override
        public void afterConnectionEstablished(WebSocketSession session) throws Exception {
            super.afterConnectionEstablished(session);
            String conversationId = session.getAttributes().get("conversationId").toString();
            IntegrationFlow f = IntegrationFlow.from(MongoDb.inboundChannelAdapter(mongoTemplate.getMongoDatabaseFactory(), 
                Query.query(Criteria.where("conversationId").is(conversationId).and("sessions").nin(session.getId())))
                .collectionName("test").entityClass(Document.class)
                    .update(new Update().push("sessions", session.getId())), 
                        p -> p.poller(pm -> pm.fixedDelay(1000L)))
            .split()
            .enrichHeaders(s -> s.header(SimpMessageHeaderAccessor.SESSION_ID_HEADER, session.getId()))
            .channel("wsMessages")
            .get();
            flowContext.registration(f).id(session.getId()).useFlowIdAsPrefix().register();
        }
        @Override
        public void afterConnectionClosed(WebSocketSession session, CloseStatus closeStatus) throws Exception {
            super.afterConnectionClosed(session, closeStatus);
            flowContext.remove(session.getId());
        }
    });
    return container;
}

@Bean
public MessageHandler webSocketOutboundAdapter(ServerWebSocketContainer container) {
    return new WebSocketOutboundMessageHandler(container);  
}

@Bean
public IntegrationFlow wsOut(MessageHandler webSocketOutboundAdapter) {
    return IntegrationFlow.from("wsMessages")
        .handle(webSocketOutboundAdapter).get();
}

@Bean 
public QueueChannelSpec wsMessages() {
    return MessageChannels.queue();
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 21:48:11