如何在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
相关产品推荐
相关产品推荐

