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

Spring非响应式SSE多客户端推送异常问题求助

Spring SSE广播问题:多客户端仅最后一个收到通知的解决方案

问题核心分析

你的代码中,emitters 使用 MutableMap<String, SseEmitter> 存储,同一个 eventId 只能关联一个 SseEmitter。当新客户端订阅同一个 eventId 时,旧的 SseEmitter 会被覆盖,导致触发事件时只有最后一个订阅的客户端能收到通知。


问题解答

1. 如何实现同时向多个客户端推送事件?

修改存储结构,让同一个 eventId 对应多个客户端的 SseEmitter,并保证线程安全:

  • 将 emitters 改为 Map<String, MutableList<SseEmitter>>,用线程安全的集合(如 ConcurrentHashMap + CopyOnWriteArrayList)避免并发修改异常
  • 客户端订阅时,将新创建的 SseEmitter 添加到对应 eventId 的列表中
  • 触发事件时,遍历对应 eventId 下的所有 SseEmitter,逐个发送事件
  • 处理 SseEmitter 的完成/超时/异常事件时,从列表中移除失效的 SseEmitter

修改后的控制器代码:

@RestController
@RequestMapping("/sse/servlet")
class EventController {
    private val objectMapper: ObjectMapper = ObjectMapper()
    private val log = KotlinLogging.logger {}
    // 线程安全的存储结构:eventId -> 对应客户端的SseEmitter列表
    private val emitters: MutableMap<String, MutableList<SseEmitter>> = ConcurrentHashMap()

    @GetMapping("/notifications")
    fun listenNotifications(@RequestParam eventId: String): SseEmitter {
        val emitter = SseEmitter(TimeUnit.MINUTES.toMillis(10))
        // 初始化eventId对应的列表(线程安全)
        emitters.computeIfAbsent(eventId) { CopyOnWriteArrayList() }.add(emitter)

        emitter.onCompletion {
            log.info("SSE connection closed for eventId: $eventId")
            emitters[eventId]?.remove(emitter)
            // 如果列表为空,清理空条目
            emitters[eventId]?.takeIf { it.isEmpty() }?.let { emitters.remove(eventId) }
        }
        emitter.onTimeout {
            log.info("SSE connection timed out for eventId: $eventId")
            emitter.complete()
        }
        emitter.onError { throwable ->
            log.error("Listen SSE exception for eventId: $eventId", throwable)
            emitters[eventId]?.remove(emitter)
        }
        return emitter
    }

    @PostMapping("/notifications")
    @ResponseStatus(ACCEPTED)
    fun fireNotification(
        @RequestParam eventId: String,
        @RequestBody notification: Notification
    ) {
        val emitterList = emitters[eventId] ?: run {
            log.info("No SSE connections found for eventId: $eventId")
            return
        }
        // 遍历所有客户端的SseEmitter发送事件
        emitterList.forEach { sseEmitter ->
            handleEmitter(notification, sseEmitter, eventId)
        }
    }

    private fun handleEmitter(
        notification: Notification,
        sseEmitter: SseEmitter,
        eventId: String
    ) = try {
        val data = objectMapper.writeValueAsString(notification)
        val sseEventBuilder = event().data(data)
        sseEmitter.send(sseEventBuilder)
    } catch (ioException: IOException) {
        log.error("Send SSE exception for eventId: $eventId", ioException)
        // 移除发送失败的SseEmitter
        emitters[eventId]?.remove(sseEmitter)
    }
}

2. 同网络下是否只会存在一个SSE连接?

不会。每个客户端(浏览器标签页、终端窗口)发起的SSE订阅都是独立的HTTP请求,会建立独立的长连接,和网络环境无关。只要是不同的请求实例,就会生成不同的 SseEmitter,对应独立的连接。


关键注意事项

  • 线程安全:必须使用线程安全的集合存储 SseEmitter,避免多客户端并发订阅/取消时出现并发修改异常
  • 资源清理:及时移除失效的 SseEmitter(完成、超时、发送失败),避免内存泄漏
  • 超时设置:合理设置 SseEmitter 的超时时间,避免无效长连接占用资源

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 19:02:50