如何用Spring WebSocket向指定用户推送产品变更消息?
问题:Spring WebSocket中关联会话与产品ID,实现向指定用户发送产品变更消息
我是Spring WebSocket新手,希望向客户端发送产品变更消息,计划按以下方式实现:
客户端代码:
var socket = new SockJS('/websocket'); var stompClient = Stomp.over(socket); stompClient.connect({}, function (frame) { stompClient.subscribe('/product/changes', function (scoredata) { // 接收产品变更消息 }); }); // 发送Ajax请求告知服务器需订阅ID为5的产品变更 sendAjaxRequest(5);
已完成的Spring配置:
@Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Override public void registerStompEndpoints(StompEndpointRegistry registry) { registry.addEndpoint("/websocket").withSockJS(); } @Override public void configureMessageBroker(MessageBrokerRegistry registry) { registry.enableSimpleBroker("/product/"); registry.setApplicationDestinationPrefixes("/app"); } }
现在需要实现以下方法,将当前WebSocket会话与产品ID关联,后续通过convertAndSendToUser向当前用户发送变更消息:
@RestController public class ProductController { @GetMapping("product-{id}") public void startSubscribe(@PathVariable("id") Long id) { // 这里需要实现关联逻辑 } }
回答
嘿,我来帮你搞定这个问题!你的核心需求是把用户的WebSocket会话和订阅的产品ID绑定起来,之后能精准推送对应产品的变更消息。咱们一步步来实现:
第一步:创建订阅注册表
首先需要一个组件来存储「会话ID - 产品ID」的关联关系,单实例场景下用ConcurrentHashMap就足够(分布式场景可以换成Redis):
@Component public class ProductSubscriptionRegistry { // key: 会话ID(这里用HttpSessionId,因为Ajax和WebSocket同属一个会话) // value: 该会话订阅的产品ID集合 private final Map<String, Set<Long>> sessionProductMap = new ConcurrentHashMap<>(); // 添加订阅关系 public void addSubscription(String sessionId, Long productId) { sessionProductMap.computeIfAbsent(sessionId, k -> new HashSet<>()).add(productId); } // 根据产品ID获取所有订阅它的会话ID public Set<String> getSessionIdsForProduct(Long productId) { return sessionProductMap.entrySet().stream() .filter(entry -> entry.getValue().contains(productId)) .map(Map.Entry::getKey) .collect(Collectors.toSet()); } // 会话断开时清理订阅关系 public void removeSession(String sessionId) { sessionProductMap.remove(sessionId); } }
第二步:配置WebSocket会话监听
我们需要在WebSocket连接建立和断开时,同步会话信息到注册表。修改你的WebSocketConfig:
@Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { private final ProductSubscriptionRegistry subscriptionRegistry; // 构造注入注册表 public WebSocketConfig(ProductSubscriptionRegistry subscriptionRegistry) { this.subscriptionRegistry = subscriptionRegistry; } @Override public void registerStompEndpoints(StompEndpointRegistry registry) { registry.addEndpoint("/websocket").withSockJS(); } @Override public void configureMessageBroker(MessageBrokerRegistry registry) { registry.enableSimpleBroker("/product/"); registry.setApplicationDestinationPrefixes("/app"); } // 添加通道拦截器,监听连接/断开事件 @Override public void configureClientInboundChannel(ChannelRegistration registration) { registration.interceptors(new ChannelInterceptor() { @Override public Message<?> preSend(Message<?> message, MessageChannel channel) { StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); if (StompCommand.CONNECT.equals(accessor.getCommand())) { // 获取HttpSessionId,存入Stomp会话属性中 String httpSessionId = accessor.getSessionAttributes().get("HTTP_SESSION_ID").toString(); accessor.getSessionAttributes().put("HTTP_SESSION_ID", httpSessionId); } else if (StompCommand.DISCONNECT.equals(accessor.getCommand())) { // 会话断开时清理注册表 String httpSessionId = accessor.getSessionAttributes().get("HTTP_SESSION_ID").toString(); subscriptionRegistry.removeSession(httpSessionId); } return message; } }); } }
第三步:实现startSubscribe方法
现在可以在ProductController里关联会话和产品ID了,同时注入SimpMessagingTemplate用来发送测试消息:
@RestController public class ProductController { private final ProductSubscriptionRegistry subscriptionRegistry; private final SimpMessagingTemplate messagingTemplate; // 构造注入依赖 public ProductController(ProductSubscriptionRegistry subscriptionRegistry, SimpMessagingTemplate messagingTemplate) { this.subscriptionRegistry = subscriptionRegistry; this.messagingTemplate = messagingTemplate; } @GetMapping("product-{id}") public void startSubscribe(@PathVariable("id") Long id, HttpSession session) { // 获取当前Ajax请求的HttpSessionId(和WebSocket会话同属一个) String sessionId = session.getId(); // 添加订阅关系 subscriptionRegistry.addSubscription(sessionId, id); // 测试:发送一条订阅成功的消息给当前用户 messagingTemplate.convertAndSendToUser(sessionId, "/product/changes", "已成功订阅产品ID: " + id); } // 示例:产品变更时触发的方法(你可以根据实际业务场景调用) public void onProductChanged(Long productId, String changeContent) { // 获取所有订阅该产品的会话ID Set<String> sessionIds = subscriptionRegistry.getSessionIdsForProduct(productId); // 给每个订阅的用户推送消息 for (String sessionId : sessionIds) { messagingTemplate.convertAndSendToUser(sessionId, "/product/changes", String.format("产品ID %d 发生变更: %s", productId, changeContent)); } } }
第四步:修改前端订阅地址
注意!使用convertAndSendToUser时,客户端需要订阅/user/product/changes(Spring会自动为用户消息加上/user/前缀),所以修改你的前端代码:
var socket = new SockJS('/websocket'); var stompClient = Stomp.over(socket); stompClient.connect({}, function (frame) { // 这里要改成/user/product/changes stompClient.subscribe('/user/product/changes', function (scoredata) { console.log('收到产品变更消息:', scoredata.body); // 处理消息的逻辑 }); }); // 发送Ajax请求告知服务器需订阅ID为5的产品变更 sendAjaxRequest(5);
额外说明
- 如果你的应用是分布式多实例部署,内存中的
ConcurrentHashMap就不适用了,需要换成Redis这样的分布式缓存来存储订阅关系,确保所有实例都能访问到相同的订阅数据。 - 如果是无状态认证(比如JWT),可以用用户名代替HttpSessionId来关联订阅关系,这样更适合前后端分离的场景。
内容的提问来源于stack exchange,提问作者Morteza Malvandi
相关产品推荐
相关产品推荐

