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

Spring Boot MVC下基于SSE实现单播消息推送至指定客户端的方案咨询

基于Spring SSE实现单播推送的实现方案

改造思路

你原有代码为每个客户端连接创建独立的SseEmitter并各自循环推送消息,属于每个连接单独发送内容的模式,要实现单播只需要建立用户唯一标识与SseEmitter的绑定映射,业务触发推送时只向目标用户对应的Emitter发送消息即可。

核心实现步骤

    1. 维护线程安全的Emitter映射表
      使用ConcurrentHashMap存储映射关系,key为客户端唯一标识(可采用用户ID、设备唯一ID等能唯一区分目标客户端的字段),value为对应的SseEmitter实例,同时配置回调自动清理失效实例避免内存泄漏。
    1. 改造SSE订阅接口
      要求客户端订阅时携带唯一标识,完成Emitter注册与回调配置
    1. 实现单播消息发送方法
      业务侧需要推送时,传入目标用户标识与消息内容,从映射表中取出对应Emitter完成发送即可

代码示例

1. 全局映射表定义

// 全局存储用户标识与SseEmitter的映射,线程安全
private static final ConcurrentHashMap<String, SseEmitter> EMITTER_MAP = new ConcurrentHashMap<>();
// 自定义超时时间,可根据业务调整,设置为0则永不超时(不推荐)
private static final Long SSE_TIMEOUT = 60000L;

2. 改造后的订阅接口

@GetMapping("/sse-emitter")
public SseEmitter sseEmitter(@RequestParam String clientId) {
    // 先清除同个clientId的旧连接,避免重复订阅
    if (EMITTER_MAP.containsKey(clientId)) {
        SseEmitter oldEmitter = EMITTER_MAP.remove(clientId);
        try {
            oldEmitter.complete();
        } catch (Exception e) {
            // 忽略旧连接关闭异常
        }
    }
    SseEmitter emitter = new SseEmitter(SSE_TIMEOUT);
    // 配置回调
    emitter.onCompletion(() -> EMITTER_MAP.remove(clientId));
    emitter.onTimeout(() -> EMITTER_MAP.remove(clientId));
    emitter.onError((throwable) -> EMITTER_MAP.remove(clientId));
    EMITTER_MAP.put(clientId, emitter);
    return emitter;
}

3. 单播消息发送方法

/**
 * 向单个客户端推送消息
 * @param clientId 目标客户端唯一标识
 * @param message 推送的消息内容
 */
public void sendSingleMessage(String clientId, Object message) {
    SseEmitter emitter = EMITTER_MAP.get(clientId);
    if (emitter == null) {
        // 客户端未订阅,可根据业务处理,比如存离线消息待客户端上线后推送
        return;
    }
    try {
        SseEmitter.SseEventBuilder event = SseEmitter.event()
                .name("SINGLE_NOTIFY")
                .data(message);
        emitter.send(event);
    } catch (Exception e) {
        // 发送失败说明连接失效,移除映射
        EMITTER_MAP.remove(clientId);
        try {
            emitter.completeWithError(e);
        } catch (Exception ex) {
            // 忽略关闭异常
        }
    }
}

注意事项

  • 移动端(Android/iOS)客户端需要实现SSE自动重连逻辑,网络波动、连接超时后需要携带相同的clientId重新发起订阅请求
  • 如果需要支持离线消息,可以在发送时发现客户端未在线,将消息存入缓存/数据库,待客户端重新订阅后先推送离线消息
  • SSE适合推送小体积的通知类消息,大体积数据建议客户端收到通知后主动调用接口拉取
  • 生产环境建议对映射表做定期巡检,清理长时间无活跃的无效连接,避免内存溢出

内容的提问来源于stack exchange,提问作者Zakir saifi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 20:15:01