如何创建可在未来发布任意事件的自定义Flux?
实现可动态发布事件的Flux流
要创建能在未来任意时间主动发布事件的Flux,Flux.from()和Flux.generate()确实不适用——前者依赖现有数据源,后者是同步生成逻辑,都无法支持外部触发的动态事件。你需要用Reactor提供的Processor + Sink组合来实现。
具体实现示例
这里以EmitterProcessor为例,它支持多播事件,还可以配置缓存历史消息,适合Web流场景:
import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.EmitterProcessor; import reactor.core.publisher.FluxSink; @RestController public class EventController { // 用于保存事件流的处理器 private final EmitterProcessor<String> eventProcessor; // 用于外部发送事件的Sink private final FluxSink<String> eventSink; public EventController() { // 创建无初始缓存的Processor,新订阅者不会收到历史消息 this.eventProcessor = EmitterProcessor.create(false); // 获取Sink,后续通过它发送事件 this.eventSink = eventProcessor.sink(); } @GetMapping(path = "/event/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> eventStream() { // 返回Processor的Flux视图,供客户端订阅事件流 return eventProcessor; } // 新增接口,用于触发发布任意事件 @PostMapping("/event/publish") public ResponseEntity<Void> publishEvent(@RequestBody String event) { // 发送事件到流中,所有订阅的客户端都会收到 eventSink.next(event); return ResponseEntity.ok().build(); } }
关键细节说明
- 线程安全:
FluxSink本身是线程安全的,可以从任意线程调用next()方法发送事件,比如定时任务、其他接口请求都能触发。 - 历史消息控制:如果希望新订阅者能收到最近的N条历史事件,可以用
EmitterProcessor.create(N),比如create(1)会保留最后一条事件。 - 轻量替代方案:如果不需要历史消息缓存,可改用
DirectProcessor,它更轻量,但新订阅者只能收到订阅后的事件:private final DirectProcessor<String> eventProcessor = DirectProcessor.create(); private final FluxSink<String> eventSink = eventProcessor.sink(); - 资源管理:如果需要在所有订阅者取消后清理资源,可以给Processor添加终止钩子:
eventProcessor.onTerminate().doFinally(signal -> { // 执行资源清理逻辑,比如关闭连接、释放资源 }).subscribe();
核心逻辑
通过Processor维护事件流的订阅关系,Sink作为外部触发事件的入口,这样就能在任意时间调用eventSink.next()来发布事件,所有订阅了/event/stream的客户端都会实时收到消息。
内容的提问来源于stack exchange,提问作者lance-java
相关产品推荐
相关产品推荐

