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
相关产品推荐
相关产品推荐

