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

