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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 03:59:07