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

首次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;
            }
    }
}
问题分析与解决方案

核心原因

  1. 请求线程被阻塞:Flux.create的accept回调默认在订阅者线程(即Web容器的请求处理线程)执行,take()是阻塞方法,首次调用时线程卡在等待队列事件,无法处理心跳或返回数据,导致请求无限阻塞。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 09:44:50