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

如何在Flux的doOnError中获取触发错误的Item对象

解决响应式流中捕获错误时获取对应Item的问题

在你当前的代码里,doOnError 无法拿到触发错误的Item,因为它是监听整个流的错误信号,此时已经丢失了具体Item的上下文。要解决这个问题,你需要把保存操作和错误处理绑定到每个Item的处理阶段,确保在能访问到Item的时候捕获错误。

以下是两种可行的实现方式:

方式一:使用 flatMap + onErrorResume

把保存操作包装成响应式类型(比如Mono),然后利用onErrorResume在出错时直接获取当前Item并处理:

reactiveTemplate.query(cql, itemMapper)
    .flatMap(item -> 
        // 如果saveRecord是同步方法,用Mono.fromRunnable包装;如果本身是异步返回Mono,直接调用即可
        Mono.fromRunnable(() -> saveRecord(item))
            .onErrorResume(th -> {
                // 这里可以直接拿到出错的item和异常
                processErrorRecord(item, th);
                // 返回Mono.empty()让流继续处理下一个Item,若要终止流则返回Mono.error(th)
                return Mono.empty();
            })
    );

如果你的saveRecord本身就是返回Mono<Void>的异步方法,代码可以简化为:

reactiveTemplate.query(cql, itemMapper)
    .flatMap(item -> saveRecord(item)
            .onErrorResume(th -> {
                processErrorRecord(item, th);
                return Mono.empty();
            })
    );

方式二:使用 handle 操作符

handle 操作符允许你逐个处理元素,同时捕获异常并自定义错误处理逻辑:

reactiveTemplate.query(cql, itemMapper)
    .handle((item, sink) -> {
        try {
            saveRecord(item);
            // 如果需要把处理后的Item传递到下游,调用sink.next(item)
            sink.next(item);
        } catch (Throwable th) {
            // 直接拿到出错的item和异常
            processErrorRecord(item, th);
            // 若不想终止整个流,就不要调用sink.error(th);否则调用该方法终止流
            // sink.error(th);
        }
    });

核心原理

原来的doOnNext属于副作用操作,它不参与流的元素传递上下文,一旦出错,错误会直接流到下游的doOnError,此时已经无法关联到具体的Item。而上面两种方式都是在每个Item的处理上下文内捕获错误,因此能直接访问到触发错误的Item。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 07:34:57