Spring Websocket+Stomp集成RabbitMQ场景下心跳监听器/拦截器实现咨询
解决方案
你可以通过以下三个层级的方案组合实现需求,完全覆盖心跳刷新、节点宕机清理的场景:
1. 直接监听Spring封装的心跳事件(最简方案)
Spring Websocket 已经原生封装了SessionHeartbeatEvent事件,不需要自行解析Stomp帧,直接监听该事件即可拿到会话ID刷新缓存:
@Component public class WebsocketHeartbeatListener implements ApplicationListener<SessionHeartbeatEvent> { private final HazelcastInstance hazelcastInstance; public WebsocketHeartbeatListener(HazelcastInstance hazelcastInstance) { this.hazelcastInstance = hazelcastInstance; } @Override public void onApplicationEvent(SessionHeartbeatEvent event) { String sessionId = event.getSessionId(); // get操作会自动重置Hazelcast Map条目的最大空闲倒计时 hazelcastInstance.getMap("websocket_sessions").get(sessionId); } }
注意:该事件默认在收到客户端心跳、服务端发送心跳时都会触发,不需要额外配置
2. 自定义Stomp帧拦截器(兼容自定义心跳规则场景)
如果你的心跳规则有定制化修改,没有触发Spring默认的心跳事件,可以通过拦截Stomp入站通道的所有帧,自行识别心跳包:
第一步:实现心跳拦截器
@Component public class StompHeartbeatInterceptor implements ChannelInterceptor { private final HazelcastInstance hazelcastInstance; public StompHeartbeatInterceptor(HazelcastInstance hazelcastInstance) { this.hazelcastInstance = hazelcastInstance; } @Override public Message<?> preSend(Message<?> message, MessageChannel channel) { StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); if (accessor == null) { return message; } String sessionId = accessor.getSessionId(); if (sessionId == null) { return message; } // 所有携带会话ID的入站帧(包括心跳、业务消息)都刷新缓存,减少残留概率 hazelcastInstance.getMap("websocket_sessions").get(sessionId); return message; } }
第二步:注册拦截器到Stomp通道
@Configuration @EnableWebSocketMessageBroker public class WebsocketConfig implements WebSocketMessageBrokerConfigurer { private final StompHeartbeatInterceptor heartbeatInterceptor; public WebsocketConfig(StompHeartbeatInterceptor heartbeatInterceptor) { this.heartbeatInterceptor = heartbeatInterceptor; } @Override public void configureClientInboundChannel(ChannelRegistration registration) { registration.interceptors(heartbeatInterceptor); } // 保留原有其他配置 }
Stomp标准心跳帧的特征是无Command指令、Payload为空,上述代码兼容所有合法Stomp帧的刷新逻辑
3. Hazelcast侧节点宕机兜底清理
针对服务节点异常宕机的场景,仅靠空闲超时会有最高等于超时时间的残留窗口,可以通过节点监听机制直接清理离线节点的所有会话:
- 存储会话缓存时,额外存入该会话归属的服务节点唯一标识
- 监听Hazelcast集群的节点离线事件,直接删除所有归属离线节点的会话:
@Component public class HazelcastNodeOfflineListener implements MembershipListener { private final HazelcastInstance hazelcastInstance; // 当前节点唯一标识,可固定配置或启动时自动生成 private final String currentNodeId = UUID.randomUUID().toString(); public HazelcastNodeOfflineListener(HazelcastInstance hazelcastInstance) { this.hazelcastInstance = hazelcastInstance; hazelcastInstance.getCluster().addMembershipListener(this); } @Override public void memberRemoved(MembershipEvent event) { String offlineNodeId = event.getMember().getUuid().toString(); IMap<String, SessionMeta> sessionMap = hazelcastInstance.getMap("websocket_sessions"); // 批量删除离线节点的所有会话 sessionMap.removeAll(Predicates.equal("ownerNodeId", offlineNodeId)); } // 新增会话时调用,存入归属节点ID public void addSession(String sessionId, SessionMeta meta) { meta.setOwnerNodeId(currentNodeId); // 最大空闲时间建议设置为客户端心跳间隔的3倍 sessionMap.put(sessionId, meta, 3, TimeUnit.MINUTES); } }
注意事项
- Hazelcast Map需要提前配置
max-idle-seconds参数,默认情况下get/contains操作都会重置条目的空闲倒计时 - 建议业务请求帧也同步触发缓存刷新,不需要完全依赖心跳,避免网络波动导致的误删
内容的提问来源于stack exchange,提问作者BlackLog
相关产品推荐
相关产品推荐

