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

