如何从其他方法发送ServerSentEvent并通过Flux消费Channel
SSE事件触发推送改造方案
你可以通过Reactor提供的Sinks组件实现事件驱动的SSE推送,替换原有的定时轮询逻辑,步骤如下:
- 构建线程安全的事件广播组件,按子域名隔离不同租户的事件流,同时支持业务侧发布事件、SSE侧订阅事件
- 改造原有SSE接口,移除
Flux.interval定时逻辑,改为订阅事件流,同时补充初始数据下发、心跳保活逻辑避免连接异常 - 在数据保存完成的业务节点,调用广播组件的推送方法,主动触发SSE消息下发
代码实现
1. 事件广播组件
import org.springframework.stereotype.Component import reactor.core.publisher.Flux import reactor.core.publisher.Sinks import java.util.concurrent.ConcurrentHashMap @Component class ImageEventPublisher { // 按子域名存储独立的事件广播器,避免不同租户消息串流 private val sinkMap = ConcurrentHashMap<String, Sinks.Many<MutableList<Image>>>() private fun getOrCreateSink(subdomain: String): Sinks.Many<MutableList<Image>> { return sinkMap.computeIfAbsent(subdomain) { // 多播模式:仅向订阅后在线的客户端推送消息,自动适配背压 Sinks.many().multicast().directBestEffort() } } /** * 数据保存完成后调用该方法推送更新 */ fun publishUpdate(subdomain: String, images: MutableList<Image>) { getOrCreateSink(subdomain).tryEmitNext(images) } /** * 提供给SSE控制器的事件订阅流 */ fun subscribe(subdomain: String): Flux<MutableList<Image>> { return getOrCreateSink(subdomain).asFlux() } }
2. 改造SSE控制器
import org.springframework.http.codec.ServerSentEvent import org.springframework.web.bind.annotation.GetMapping import jakarta.servlet.http.HttpServletRequest import reactor.core.publisher.Flux import java.time.Duration @GetMapping("/images-sse") fun getImagesAsSSE( request: HttpServletRequest ): Flux<ServerSentEvent<MutableList<Image>>> { val subdomain = request.serverName.split(".").first() // 连接建立后先推送一次全量数据,避免客户端等待 val initialData = Flux.just( ServerSentEvent.builder<MutableList<Image>>() .event("image-update") .data(weddingService.getBySubdomain(subdomain)?.pictures ?: mutableListOf()) .build() ) // 订阅后续数据更新事件 val updateStream = imageEventPublisher.subscribe(subdomain) .map { images -> ServerSentEvent.builder<MutableList<Image>>() .event("image-update") .data(images) .build() } // 合并初始数据、更新流,加30秒间隔心跳防止连接被网关/浏览器超时断开 return Flux.merge(initialData, updateStream) .onBackpressureBuffer() .mergeWith( Flux.interval(Duration.ofSeconds(30)) .map { ServerSentEvent.builder<MutableList<Image>>().comment("heartbeat").build() } ) }
3. 业务保存逻辑触发推送
在你完成图片数据保存的代码位置,注入ImageEventPublisher,保存成功后调用推送方法即可:
// 原有图片/数据保存逻辑执行完成后 val latestImages = weddingService.getBySubdomain(subdomain)?.pictures ?: mutableListOf() imageEventPublisher.publishUpdate(subdomain, latestImages)
注意:以上实现基于应用内存存储事件流,仅适用于单实例部署场景。如果是多实例部署,需要额外引入Redis Pub/Sub、消息队列等组件实现跨实例的事件同步。
内容的提问来源于stack exchange,提问作者tschuehly
相关产品推荐
相关产品推荐

