Flux.fromStream抛出异常如何处理?如何忽略S3Exception继续执行后续流程
解决方案
核心思路是把异常处理逻辑放在每个独立的fetchResources调用链中,避免单条请求异常终止整个响应流:
具体修改点
- 移除全局的
doOnError配置,将错误捕获逻辑下沉到flatMap内部的fetchResources流上 - 使用
onErrorResume操作符捕获单条请求的异常,打印日志后返回Mono.empty(),这样出错的请求不会发射任何元素,也不会中断整体流程 - 你原代码里的
.flatMap(idS3Object -> Mono.just(s3Object))属于笔误,示例代码里做了对应修正,如果实际逻辑不同可自行调整
修改后代码
public Mono<ResponseEntity<Data>> getData(@RequestParam List<String> tagIds){ Flux<S3Object> s3ObjectFlux = Flux.fromStream(tagIds.stream()) .parallel() .runOn(Schedulers.boundedElastic()) // 错误处理放在每个单独的fetch请求内部 .flatMap(id -> fetchResources(id) // 捕获当前请求的所有异常 .onErrorResume(throwable -> { log.error("拉取tagId:{}对应的S3资源失败", id, throwable); return Mono.empty(); }) ) .ordered((u1, u2) -> u2.hashCode() - u1.hashCode()); Mono<Data> data = s3ObjectFlux.collectList() .map(Data::new); // 补全返回值逻辑,原代码漏了返回ResponseEntity的部分 return data.map(ResponseEntity::ok); }
原理解释
onErrorResume会在当前流发生异常时切换到备用流,这里返回Mono.empty()就等于直接丢弃这条出错的请求结果,不会影响其他并行请求的处理- 所有请求处理完成后,
collectList会把正常返回的S3对象收集成列表,哪怕所有请求都失败,也会收集到空列表,不会抛出500错误
内容的提问来源于stack exchange,提问作者plzdontkillme
相关产品推荐
相关产品推荐

