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

使用directBestEffort的Sinks出现FAIL_OVERFLOW/FAIL_TERMINATED的原因及复现

问题分析与解决方案

错误原因解析

  • FAIL_OVERFLOW:当调用sink发送数据时,所有订阅该sink的Flux订阅者都无法处理这条数据,就会返回该结果。比如所有订阅者都处于背压状态(处理速度跟不上,自身缓冲区已满),或者还没完成订阅初始化(未调用request()请求数据),没有任何一个订阅者能接收当前数据时,就会触发这个错误。
  • FAIL_TERMINATED:这个错误说明sink已经处于终止状态(比如调用过sink.complete()/sink.error(),或者所有订阅者都已终止导致sink自动终止),此时再尝试发送数据就会返回该结果。你日志里先出现FAIL_OVERFLOW再出现FAIL_TERMINATED,大概率是某次发送失败后,sink因某种原因被终止了,后续发送操作自然失败。

文档语句解读

文档里的none of the subscribers can process an element就是字面意思:当前所有订阅该sink的订阅者,没有一个能够接收并处理当前要发送的元素。
对比正常情况:如果有至少一个订阅者能处理数据,directBestEffort会把数据发给这个快的订阅者,慢订阅者会直接丢数据,但不会触发FAIL_OVERFLOW;只有当所有订阅者都“接不住”数据时,才会返回溢出错误。

本地复现步骤

复现FAIL_OVERFLOW

// 创建directBestEffort类型的sink
Sinks.Many<String> sink = Sinks.many().multicast().directBestEffort();
Flux<String> flux = sink.asFlux();

// 添加两个慢订阅者,模拟所有订阅者都无法处理数据
flux.subscribe(
    data -> {
        // 模拟极慢的处理速度,让订阅者进入背压状态
        try { Thread.sleep(1000); } catch (InterruptedException e) {}
        System.out.println("Subscriber 1 received: " + data);
    },
    error -> System.err.println("Subscriber 1 error: " + error)
);

flux.subscribe(
    data -> {
        try { Thread.sleep(1000); } catch (InterruptedException e) {}
        System.out.println("Subscriber 2 received: " + data);
    },
    error -> System.err.println("Subscriber 2 error: " + error)
);

// 快速发送大量数据,触发所有订阅者背压
for (int i = 0; i < 100; i++) {
    Sinks.EmitResult result = sink.tryEmitNext("Data " + i);
    if (result.isFailure()) {
        System.err.println("Emit failed: " + result);
        // 这里会打印FAIL_OVERFLOW
    }
}

复现FAIL_TERMINATED

在上面的代码基础上,触发FAIL_OVERFLOW后,手动终止sink再继续发送数据:

// 终止sink
sink.tryEmitComplete();

// 再发送数据,就会返回FAIL_TERMINATED
Sinks.EmitResult result = sink.tryEmitNext("After termination data");
if (result.isFailure()) {
    System.err.println("Emit failed: " + result); // 输出FAIL_TERMINATED
}

或者模拟订阅者全部终止的场景:

// 创建订阅者并立即取消
Disposable sub1 = flux.subscribe(data -> {});
Disposable sub2 = flux.subscribe(data -> {});
sub1.dispose();
sub2.dispose();

// 此时所有订阅者都已终止,sink会进入终止状态,发送数据就会返回FAIL_TERMINATED
Sinks.EmitResult result = sink.tryEmitNext("No subscribers data");
System.err.println("Emit failed: " + result);

额外注意点

你用到的onErrorContinue()只会处理Flux订阅者端的错误,不会影响sink的emit结果判断。如果要处理sink的发送失败,需要在调用sink.emitNext()/sink.tryEmitNext()时主动判断返回的Sinks.EmitResult,根据结果做相应处理(比如降级、丢弃或重试)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 16:22:39