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

Flux.fromStream迭代时单个元素抛异常后如何继续处理剩余元素

解决方案

核心问题原因

你之前将onError*系列操作符作用于整个Flux流,只要任意一个元素触发异常,就会抛出全局错误信号,终止整个流的处理。要实现单个元素出错不影响整体流程,需要将错误处理下沉到每个元素的独立处理逻辑中。

修改后的实现代码

public Mono<ResponseEntity<Data>> getData(@RequestParam List<String> tagIds){
   
   Flux<S3Object> s3ObjectFlux = Flux.fromStream(tagIds.stream())
           .parallel()
           .runOn(Schedulers.boundedElastic())
           // 错误处理内联到每个id的拉取逻辑中
           .flatMap(id -> fetchResources(id)
                   // 捕获当前id拉取的所有异常
                   .onErrorResume(e -> {
                       // 打印异常日志,可补充tagId等上下文信息方便排查
                       log.error("拉取tagId[{}]对应的S3资源失败", id, e);
                       // 返回空Mono,相当于过滤掉当前出错元素,不会中断整个流
                       return Mono.empty();
                   })
           )
       .sequential() 
           .ordered((u1, u2) -> u2.hashCode() - u1.hashCode());

   Mono<Data> data = s3ObjectFlux.collectList()
           .map(s3Objects -> new Data(s3Objects));
   // 补充返回ResponseEntity的逻辑即可,比如返回ResponseEntity.ok(data)
}

补充说明

如果你的fetchResources是同步方法,调用时会直接抛出异常而非返回错误的Mono,需要先用Mono.fromCallable包裹方法调用,再做错误处理:

.flatMap(id -> Mono.fromCallable(() -> fetchResources(id))
        .onErrorResume(e -> {
            log.error("拉取tagId[{}]对应的S3资源失败", id, e);
            return Mono.empty();
        })
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 18:06:07