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

如何从其他方法发送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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 16:36:29