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推送给客户端
目前代码在本地可以正常运行,但存在两个疑问:
- 是否有更简洁自然的实现方式?
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
相关产品推荐
相关产品推荐

