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

基于Project Reactor与Spring WebFlux实现异步事件推送的求助

解决方案

核心用Spring Reactor的Sinks实现异步事件的发布-订阅,确保新连接的客户端只接收订阅后的新事件,不会拿到历史消息。

具体实现代码

import reactor.core.publisher.Sinks;
import reactor.core.publisher.Flux;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.context.event.EventListener;
import org.springframework.scheduling.annotation.Async;
import org.springframework.http.MediaType;

@RestController
public class ModelController {
    // 多播模式:多个订阅者共享事件流,新订阅者仅接收订阅后的事件
    private final Sinks.Many<MyModel> modelSink = Sinks.many().multicast().onBackpressureBuffer();

    @EventListener
    @Async
    protected void handleModelChanges(ModelChangeEvent e) {
        log.debug("Event received {}", e);
        // 把事件转换为业务需要的MyModel对象
        MyModel model = convertEventToModel(e);
        // 非阻塞推送事件,避免发布线程被阻塞
        modelSink.tryEmitNext(model);
    }

    @GetMapping(value = "/model-changes", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<@NonNull MyModel> modelChanges() {
        // 返回Sink对应的Flux,供客户端订阅
        return modelSink.asFlux();
    }

    // 自定义事件转Model的逻辑,根据你的业务实现
    private MyModel convertEventToModel(ModelChangeEvent e) {
        return new MyModel(...); // 填充业务字段
    }
}

关键细节说明

  • Sinks.many().multicast().onBackpressureBuffer():这个配置是核心,既支持多客户端同时订阅,又保证新订阅者不会拿到订阅前的历史事件,同时用缓冲区处理突发的事件流量。
  • tryEmitNext:非阻塞式推送,若订阅者处理速度跟不上,会把事件暂存到缓冲区,也可以根据业务需求换成emitNext并自定义错误处理逻辑。
  • produces = MediaType.TEXT_EVENT_STREAM_VALUE:明确告诉客户端这是Server-Sent Events(SSE)流,WebClient可以直接订阅这个持续的事件流。

WebClient订阅示例

如果用WebClient对接这个端点,代码示例如下:

webClient.get()
        .uri("/model-changes")
        .accept(MediaType.TEXT_EVENT_STREAM)
        .retrieve()
        .bodyToFlux(MyModel.class)
        .subscribe(model -> {
            // 处理新收到的MyModel事件
            log.info("Got new model: {}", model);
        });

注意事项

  • 要确保Spring开启了异步支持,需要在配置类上添加@EnableAsync注解,否则@Async不会生效。
  • 如果客户端断开重连,重连后会从断开后的新事件开始接收,符合你不接收旧事件的需求。
  • 可以根据业务场景处理tryEmitNext的返回值(Sinks.EmitResult),避免事件推送失败时静默无提示。

内容的提问来源于stack exchange,提问作者fer.marino

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 16:57:05