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

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侧节点宕机兜底清理

针对服务节点异常宕机的场景,仅靠空闲超时会有最高等于超时时间的残留窗口,可以通过节点监听机制直接清理离线节点的所有会话:

  1. 存储会话缓存时,额外存入该会话归属的服务节点唯一标识
  2. 监听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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 06:54:01