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

Quarkus/Kotlin中SseBroadcaster遇已关闭EventSink广播抛空指针

解决Quarkus Kotlin SSE广播中客户端断开引发的空指针与连接清理问题

问题根源

JAX-RS默认的SseBroadcaster存在两个核心问题:

  1. 同步广播机制:broadcast()方法发送消息时,若某一SseEventSink已失效(如客户端主动断开),发送异常会直接中断后续所有客户端的消息推送
  2. 无主动清理能力:广播器未提供移除失效Sink的API,仅通过监听器通知连接状态变化,无法自动清理无效连接

解决方案

放弃依赖默认SseBroadcaster,改为手动维护线程安全的活跃SseEventSink集合,实现以下逻辑:

  • 用线程安全集合管理所有活跃连接
  • 为每个Sink注册关闭/错误回调,主动移除失效实例
  • 广播时逐个处理Sink发送逻辑,捕获异常并隔离失效连接,避免全局中断

修改后的代码示例

【Service】

package your.packagename.services

import jakarta.enterprise.context.ApplicationScoped
import jakarta.ws.rs.core.Context
import jakarta.ws.rs.sse.OutboundSseEvent
import jakarta.ws.rs.sse.Sse
import jakarta.ws.rs.sse.SseEventSink
import java.util.concurrent.CopyOnWriteArrayList

@ApplicationScoped
class EventService(@Context private val sse: Sse) {

    // 线程安全集合,适配SSE广播读多写少的场景
    private val activeSinks = CopyOnWriteArrayList<SseEventSink>()

    fun createEvent(message: String): OutboundSseEvent {
        return sse.newEvent(message)
    }

    fun register(sseEventSink: SseEventSink) {
        println("Added new event sink $sseEventSink")
        
        // 连接关闭时主动从集合移除并释放资源
        sseEventSink.onClose {
            println("Closing event stream: $sseEventSink")
            activeSinks.remove(sseEventSink)
            if (!sseEventSink.isClosed) {
                sseEventSink.close()
            }
        }
        
        // 发生异常时主动清理失效连接
        sseEventSink.onError { throwable ->
            println("Error on event stream ${throwable.message}: $sseEventSink")
            activeSinks.remove(sseEventSink)
            if (!sseEventSink.isClosed) {
                sseEventSink.close()
            }
        }
        
        activeSinks.add(sseEventSink)
    }

    fun broadcast(sseEvent: OutboundSseEvent) {
        println("Publishing new message ${sseEvent.data}")
        
        // 遍历所有活跃连接,单个连接失败不影响全局广播
        activeSinks.forEach { sink ->
            try {
                if (!sink.isClosed) {
                    sink.send(sseEvent)
                } else {
                    activeSinks.remove(sink)
                }
            } catch (e: Exception) {
                println("Failed to send to sink $sink: ${e.message}")
                activeSinks.remove(sink)
                if (!sink.isClosed) {
                    sink.close()
                }
            }
        }
    }
}

关键修改说明

  1. 线程安全集合:使用CopyOnWriteArrayList保证多线程环境下的并发安全,适合SSE广播场景下读操作远多于写操作的特性
  2. 连接生命周期管理:每个Sink注册独立的onClose和onError回调,主动从集合中移除失效连接,避免无效资源占用
  3. 隔离式广播:遍历集合逐个发送消息,捕获所有发送异常,单个连接的错误不会中断其他客户端的消息推送
  4. 资源泄漏防护:每次操作后检查Sink状态,确保失效连接被及时关闭并移除

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 10:40:58