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

使用flatMapDelayError时Flux遇异常后积压任务未正确执行问题咨询

问题原因分析

核心问题出在flatMapDelayError的特性和groupBy的错误处理逻辑的冲突上:

  • flatMapDelayError的作用是延迟错误到所有内部Publisher执行完成后再向下游合并抛出异常,而不是让错误不会流向下游
  • 当第一个flatMapDelayError处理完所有元素后,会把收集到的0对应的IllegalStateException向下游发射给groupBy
  • groupBy操作符只要收到上游的onError信号,无论之前已经产生的分组还有多少元素没处理,都会立刻终止所有分组流并转发错误,导致分组后的collectList还没收集完该分组的全部正常元素就提前报错,所以归约逻辑完全没有执行的机会
修复方案

要实现「先处理完所有正常数据,最后再抛出异常」的需求,我们需要先把错误转换成普通流元素,避免流提前触发onError终止,等所有正常数据处理完成后,再统一抛出收集到的错误:

修改后代码

public class Example1FlatMapDelayError1 {
    // 定义容器类区分正常数据和错误
    private static class ResultHolder {
        private final Integer value;
        private final Throwable error;

        private ResultHolder(Integer value, Throwable error) {
            this.value = value;
            this.error = error;
        }

        public static ResultHolder success(Integer value) {
            return new ResultHolder(value, null);
        }

        public static ResultHolder error(Throwable error) {
            return new ResultHolder(null, error);
        }

        public boolean isSuccess() {
            return error == null;
        }

        public Integer getValue() {
            return value;
        }

        public Throwable getError() {
            return error;
        }
    }

    public static void main(String[] args) {
        ReactorDebugAgent.init();
        ReactorDebugAgent.processExistingClasses();

        Flux.just(1, 2, 3, 4, 5, 6, 7, 8, 9, 0)
            // 先把所有元素转成普通容器,不触发onError
            .map(integer -> {
                if (integer == 0) {
                    return ResultHolder.error(new IllegalStateException());
                }
                return ResultHolder.success(integer);
            })
            // 分流:收集错误,处理正常数据
            .publish(flux -> {
                // 收集所有错误
                Flux<Throwable> errors = flux.filter(holder -> !holder.isSuccess())
                    .map(ResultHolder::getError);
                // 处理正常数据的归约逻辑
                Flux<Integer> processResults = flux.filter(ResultHolder::isSuccess)
                    .map(ResultHolder::getValue)
                    .groupBy(integer -> integer % 2 == 0 ? "Even" : "Odd")
                    .flatMap(group -> {
                        System.out.println("Grouping Key-->" + group.key());
                        return group
                            .collectList()
                            .doOnNext(integers -> System.out.println("Key Value -->" + integers.stream().map(String::valueOf).collect(Collectors.joining(","))))
                            .flatMap(integers -> Flux.fromIterable(integers).reduce((i1, i2) -> i1 * i2));
                    });
                // 先返回所有处理结果,最后合并错误抛出
                return processResults.concatWith(errors.flatMap(Mono::error));
            })
            .subscribe(
                out -> System.out.println("Final -> " + out),
                throwable -> throwable.printStackTrace()
            );
    }
}

执行结果说明

修改后代码会先输出:

Grouping Key-->Odd
Key Value -->1,3,5,7,9
Final -> 945
Grouping Key-->Even
Key Value -->2,4,6,8
Final -> 384

等所有正常数据归约完成后,再抛出IllegalStateException异常,完全符合预期。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 16:27:01