Spring Websocket STOMP订阅topic时如何在连接后获取最新历史消息?
解决方案
注意:Spring 自带的 Simple Broker 是内存级消息代理,本身不支持历史消息留存和回溯,你需要自行在业务层实现最新消息的持久化存储(比如存在Redis、数据库或者内存变量中),以下两种方案都依赖你自行实现的获取最新历史数据的逻辑。
方案1:客户端订阅后主动拉取(实现简单,推荐)
该方案无需修改WebSocket核心配置,客户端订阅完成后主动发起请求拉取最新状态即可。
服务端修改
新增消息接口用于返回最新历史数据:
@Controller public class SenhaMessageController { // 自行注入你实现的历史消息存储服务 @Autowired private SenhaMessageService senhaMessageService; @Autowired private SimpMessagingTemplate messagingTemplate; @MessageMapping("/latest-senha") public void pullLatestSenha(StompHeaderAccessor accessor) { Object latestMessage = senhaMessageService.getLatestSenhaMessage(); if (latestMessage == null) { return; } // 仅向当前请求的客户端推送最新数据,不会广播给所有订阅者 messagingTemplate.convertAndSendToUser( accessor.getSessionId(), "/topic/senha", latestMessage, Collections.singletonMap("msgType", "history") ); } }
客户端修改
在afterConnected方法中新增主动拉取逻辑:
@Override public void afterConnected(final StompSession session, final StompHeaders connectedHeaders) { // 原有订阅逻辑保持不变 session.subscribe("/topic/senha", this); // 新增:主动请求最新历史消息 session.send("/app/latest-senha", null); }
拉取到的历史消息和后续实时推送的消息都会进入handleFrame方法处理,你可以通过自定义的msgType头信息区分历史消息和实时消息。
方案2:服务端拦截订阅事件主动推送
该方案客户端无需修改任何逻辑,服务端监听到新的订阅请求时主动推送最新消息给新连接的客户端。
服务端修改
- 实现订阅事件拦截器:
@Component public class SenhaSubscribeInterceptor implements ChannelInterceptor { @Autowired private SenhaMessageService senhaMessageService; @Autowired private SimpMessagingTemplate messagingTemplate; @Override public Message<?> preSend(Message<?> message, MessageChannel channel) { StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); // 仅处理订阅请求 if (!StompCommand.SUBSCRIBE.equals(accessor.getCommand())) { return message; } // 仅处理目标topic的订阅请求 String destination = accessor.getDestination(); if (!"/topic/senha".equals(destination)) { return message; } // 推送最新历史消息给当前订阅的客户端 Object latestMessage = senhaMessageService.getLatestSenhaMessage(); if (latestMessage != null) { messagingTemplate.convertAndSendToUser( accessor.getSessionId(), "/topic/senha", latestMessage, Collections.singletonMap("msgType", "history") ); } return message; } }
- 在WebSocket配置类中注册拦截器:
@Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Autowired private SenhaSubscribeInterceptor senhaSubscribeInterceptor; @Override public void configureMessageBroker(final MessageBrokerRegistry config) { config.enableSimpleBroker("/topic"); config.setApplicationDestinationPrefixes("/app"); } @Override public void registerStompEndpoints(final StompEndpointRegistry registry) { registry.addEndpoint("/ws"); } // 新增:注册订阅拦截器 @Override public void configureClientInboundChannel(ChannelRegistration registration) { registration.interceptors(senhaSubscribeInterceptor); } }
内容的提问来源于stack exchange,提问作者Thiago Sayão
相关产品推荐
相关产品推荐

