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

使用Reactor的mergeDelayError时如何捕获全部错误而非仅首个错误?

Reactor mergeDelayError 获取全部错误的实现方法

可以获取全部错误,但不能依赖mergeDelayError直接抛出所有错误——因为Reactor的流遵循**“错误即终止”**的核心语义:一旦错误被抛入主流,整个流就会停止处理后续元素或错误。你的示例代码中,第一个错误触发后,collectList因流错误终止,所以后续的错误没有机会被捕获。

问题分析

你的示例里,Flux.fromStream(Stream.of(1,2,3)).flatMap(i -> Mono.error(...))会生成三个错误,mergeDelayError确实会延迟错误直到正常元素(1、2)发送完毕,但第一个错误进入主流程后,doOnError捕获到它,随后collectList因为流错误终止,后续的错误被直接丢弃。

解决方案:包装错误为流元素

要收集所有错误,需要提前捕获每个子流的错误,将其转换为流中的普通元素(而非抛出到主流),这样流不会因错误终止,就能收集到所有错误。

步骤1:定义错误封装类

先创建一个简单的类,用来区分成功结果和错误信息:

class Result<T> {
    private final T value;
    private final Throwable error;

    private Result(T value, Throwable error) {
        this.value = value;
        this.error = error;
    }

    // 静态工厂方法
    public static <T> Result<T> success(T value) {
        return new Result<>(value, null);
    }

    public static <T> Result<T> failure(Throwable error) {
        return new Result<>(null, error);
    }

    // Getter方法
    public T getValue() { return value; }
    public Throwable getError() { return error; }
    public boolean isSuccess() { return error == null; }
}

步骤2:修改流处理逻辑

将每个子流的错误用onErrorResume捕获,包装成Result对象,确保错误不会进入主流导致终止:

Flux.mergeDelayError(10,
        // 处理会抛出错误的流,捕获错误并包装
        Flux.fromStream(Stream.of(1, 2, 3))
                .flatMap(i -> Mono.error(new Exception(i.toString()))
                        .onErrorResume(e -> Mono.just(Result.failure(e)))),
        // 处理正常流,包装为成功结果
        Mono.just(1).map(Result::success),
        Mono.just(2).map(Result::success))
        .collectList() // 收集所有结果(成功+错误)
        .block()
        .forEach(result -> {
            if (result.isSuccess()) {
                System.out.println("正常元素: " + result.getValue());
            } else {
                System.err.println("错误信息: " + result.getError().getMessage());
            }
        });

运行这段代码后,会输出所有正常元素和三个错误信息,不会因第一个错误终止流。

简化方案:仅收集错误

如果不需要保留正常元素,只想收集所有错误,可以直接忽略正常元素,只保留错误:

Flux.mergeDelayError(10,
        // 捕获错误并保留
        Flux.fromStream(Stream.of(1, 2, 3))
                .flatMap(i -> Mono.error(new Exception(i.toString()))
                        .onErrorResume(Mono::just)),
        // 忽略正常元素
        Mono.just(1).flatMap(v -> Mono.empty()),
        Mono.just(2).flatMap(v -> Mono.empty()))
        .collectList()
        .block()
        .forEach(error -> System.err.println("错误信息: " + error.getMessage()));

核心总结

  • mergeDelayError的作用是延迟错误到正常元素处理完毕,但不会改变Reactor流“错误终止”的语义;
  • 要获取全部错误,必须在子流层面提前捕获错误,将其转换为普通元素,避免错误进入主流导致整个流终止。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 17:07:50