首次API调用时为Flux Publisher添加消费者的方案咨询
问题描述
使用Spring WebFlux实现的SSE接口/bites/events返回Flux<BiteEventsResponse>,但存在以下异常:
- 应用启动后首次调用API时,请求无限阻塞,无任何事件返回
- 第二次调用API时,事件才会正常发射
已知Flux.create(eventsPublisher).share()依赖活跃消费者触发订阅逻辑,首次调用时无消费者导致生产者阻塞,需要解决消费者状态检查及历史事件处理问题。
相关代码
BiteEventsController.java
public class BiteEventsController { private final Flux<BiteEvent> biteEvent; private final int INTERVAL = 10; private final Flux<BiteEventsResponse> keepAlive = Flux.interval(Duration.ofSeconds(INTERVAL)) .map(e -> new BiteEventsResponse("iam-alive", null)); public BiteEventsController(BiteEventsPublisher eventsPublisher) { this.biteEvent = Flux.create(eventsPublisher).share(); } @GetMapping(value = "/bites/events", produces = "text/event-stream;charset=UTF-8") public Flux<BiteEventsResponse> biteEvents(HttpServletRequest request) { try { Flux<BiteEventsResponse> biteEvent = this.biteEvent.filter( event -> (event.getOwnerId() != null )) .map(event -> { return new BiteEventsResponse("message", (BiteResponse) event.getSource()); }); return Flux.merge(keepAlive, biteEvent); } catch (Exception ex) { log.error("{}", ex); return Flux.empty(); } } }
BiteEventsPublisher.java
public class BiteEventsPublisher implements Consumer<FluxSink<BiteEvent>> { @Autowired BiteEventsPublisherConfig biteEventConfig; @EventListener public void onBiteEvent(BiteEvent event) { biteEventConfig.biteEventBlockingQueue().offer(event); } @Override public void accept(FluxSink<BiteEvent> sink) { while (true) try { BiteEvent event = biteEventConfig.biteEventBlockingQueue().take(); sink.next(event); } catch (Exception ex) { log.error("Bite events publisher failed {}", ex); break; } } }
问题分析与解决方案
核心原因
- 请求线程被阻塞:
Flux.create的accept回调默认在订阅者线程(即Web容器的请求处理线程)执行,take()是阻塞方法,首次调用时线程卡在等待队列事件,无法处理心跳或返回数据,导致请求无限阻塞。 share()的触发逻辑:share()仅在第一个订阅者出现时才启动生产者,首次调用时生产者刚启动就阻塞,直到有事件入队或二次订阅。
具体修复方案
1. 后台线程执行生产者逻辑
给Flux指定独立调度器,让阻塞的生产者逻辑在后台线程运行,避免占用请求线程:
public BiteEventsController(BiteEventsPublisher eventsPublisher) { this.biteEvent = Flux.create(eventsPublisher) .subscribeOn(Schedulers.boundedElastic()) // 适配阻塞逻辑的调度器 .share(); }
修改后,请求线程不会被卡死,keepAlive的心跳能正常返回,首次调用至少能收到心跳信号,不会无限阻塞。
2. 缓存历史事件(可选)
如果需要新订阅者获取订阅前的历史事件,用ReplayProcessor替代share(),它会缓存指定数量的历史事件:
private final Flux<BiteEvent> biteEvent; public BiteEventsController(BiteEventsPublisher eventsPublisher) { // 缓存最近100条事件,可根据需求调整数量 ReplayProcessor<BiteEvent> processor = ReplayProcessor.create(100); Flux.create(eventsPublisher) .subscribeOn(Schedulers.boundedElastic()) .subscribe(processor); this.biteEvent = processor; }
无论首次还是后续订阅,都能收到缓存的历史事件,解决首次调用无事件的问题。
3. 检查活跃消费者状态
在生产者逻辑中通过FluxSink判断是否有活跃消费者,避免无消费者时的无效阻塞:
@Override public void accept(FluxSink<BiteEvent> sink) { // 订阅取消时触发的回调 sink.onCancel(() -> log.info("No active consumers, stopping event publisher")); // 仅当有活跃消费者时继续处理事件 while (!sink.isCancelled()) { try { // 带超时的poll,避免永久阻塞 BiteEvent event = biteEventConfig.biteEventBlockingQueue().poll(Duration.ofSeconds(1)); if (event != null) { sink.next(event); } } catch (Exception ex) { log.error("Bite events publisher failed {}", ex); break; } } }
用poll(Duration)替代take(),避免无事件时永久阻塞,同时通过sink.isCancelled()判断消费者状态,无消费者时退出循环节省资源。
4. 独立心跳逻辑(可选)
让每个订阅者拥有独立的心跳实例,避免多订阅者共享心跳导致的异常:
private Flux<BiteEventsResponse> keepAlive() { return Flux.interval(Duration.ofSeconds(INTERVAL)) .map(e -> new BiteEventsResponse("iam-alive", null)); } @GetMapping(value = "/bites/events", produces = "text/event-stream;charset=UTF-8") public Flux<BiteEventsResponse> biteEvents(HttpServletRequest request) { try { Flux<BiteEventsResponse> biteEvent = this.biteEvent.filter( event -> (event.getOwnerId() != null )) .map(event -> new BiteEventsResponse("message", (BiteResponse) event.getSource())); return Flux.merge(keepAlive(), biteEvent); } catch (Exception ex) { log.error("{}", ex); return Flux.empty(); } }
内容的提问来源于stack exchange,提问作者JDGuide
相关产品推荐
相关产品推荐

