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

Reactor 3.4中Sink终止后如何重新订阅及状态检测?

解决方案

一、发送前判断Sink状态

Reactor的Sink.Many提供两种核心方式提前判断状态,避免无效发送:

  1. 使用tryEmitNext()获取发送结果
    tryEmitNext()返回Sink.EmitResult枚举,可直接判断是否能安全发送:
public void emitNext(Event event){
    log(event);
    Sink.EmitResult result = unhandledEvents.tryEmitNext(event);
    if (result == Sink.EmitResult.FAIL_OVERFLOW) {
        // 处理溢出逻辑,如临时缓存或触发重建
        handleOverflow(event);
    } else if (result == Sink.EmitResult.FAIL_TERMINATED) {
        // Sink已终止,先重建再重试发送
        rebuildSink();
        unhandledEvents.tryEmitNext(event);
    }
    // 其他结果如SUCCESS、FAIL_NON_SERIALIZED可按需处理
}
  1. 直接检查Sink当前状态
    Reactor 3.4+支持通过currentState()获取Sink状态,判断是否处于ACTIVE状态:
if (unhandledEvents.currentState() == Sinks.State.ACTIVE) {
    unhandledEvents.emitNext(event, emitNextFailureHandler);
} else {
    // Sink非活跃,触发重建逻辑
    rebuildSink();
}

二、优雅重建已终止的Sink

Sink进入TERMINATED状态后无法恢复,必须重建新的Sink,并确保下游能自动切换到新数据流:

方案1:可替换Sink + Flux.defer()

将Sink改为可替换的volatile变量,结合Flux.defer()确保下游每次订阅都获取最新的Sink数据流:

class EventEmitter {
    private volatile Sinks.Many<Event> unhandledEvents = Sinks.many().multicast().onBackpressureBuffer();

    public void emitNext(Event event){
        log(event);
        Sink.EmitResult result = unhandledEvents.tryEmitNext(event);
        if (result == Sink.EmitResult.FAIL_OVERFLOW || result == Sink.EmitResult.FAIL_TERMINATED) {
            rebuildSink();
            unhandledEvents.tryEmitNext(event);
        }
    }

    public Flux<Event> getEvents() {
        // 每次订阅都获取当前最新的Sink对应的Flux
        return Flux.defer(() -> unhandledEvents.asFlux());
    }

    private void rebuildSink() {
        // 原子替换为新的Sink
        this.unhandledEvents = Sinks.many().multicast().onBackpressureBuffer();
    }
}

下游说明:如果是长期订阅,Sink重建后需重新调用subscribe();如果是动态订阅(如每次请求触发订阅),defer()会自动使用新的Sink。

方案2:自动切换Sink(无需下游手动重订阅)

通过switchOnNext()实现下游自动切换到新Sink的数据流,无需手动重新订阅:

class EventEmitter {
    private final Sinks.Many<Sinks.Many<Event>> sinkSwitcher = Sinks.many().unicast().onBackpressureBuffer();
    private volatile Sinks.Many<Event> currentSink;

    public EventEmitter() {
        // 初始化第一个Sink
        currentSink = Sinks.many().multicast().onBackpressureBuffer();
        sinkSwitcher.tryEmitNext(currentSink);
    }

    public void emitNext(Event event){
        log(event);
        Sink.EmitResult result = currentSink.tryEmitNext(event);
        if (result == Sink.EmitResult.FAIL_OVERFLOW || result == Sink.EmitResult.FAIL_TERMINATED) {
            rebuildSink();
            currentSink.tryEmitNext(event);
        }
    }

    public Flux<Event> getEvents() {
        // 自动切换到最新的Sink数据流
        return sinkSwitcher.asFlux().switchOnNext(Sinks.Many::asFlux);
    }

    private void rebuildSink() {
        Sinks.Many<Event> newSink = Sinks.many().multicast().onBackpressureBuffer();
        this.currentSink = newSink;
        sinkSwitcher.tryEmitNext(newSink);
    }
}

优势:下游只需订阅一次getEvents(),Sink重建时会自动切换到新数据流。

三、避免Sink终止的前置优化

从根源上避免因溢出导致Sink终止,可调整Sink的背压策略:

  • 增大队列容量并设置溢出丢弃策略(丢弃老元素):
private final Sinks.Many<Event> unhandledEvents = Sinks.many().multicast().onBackpressureBuffer(1024, true);

队列满时会丢弃最早的元素,而非触发FAIL_OVERFLOW导致Sink终止。


内容的提问来源于stack exchange,提问作者Dans Merino

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 06:37:08