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

如何创建可在未来发布任意事件的自定义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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 05:42:46