Spring Webflux WebClient嵌套调用问题:超时与批量完成通知需求
WebClient异步请求优化:解决超时问题+实现全量完成通知
问题背景
需要先获取项目ID列表,再逐个调用接口拉取完整项目数据存入数据库。当前实现能运行但存在以下问题:
- 偶尔出现
Operation timed out错误 - 无法统一监听所有项目的获取+保存操作完成状态
- 代码结构不规范,嵌套订阅导致异步流程失控
现有问题代码
this.webClient.get() .uri(locationURI) .accept(MediaType.APPLICATION_JSON) .exchangeToMono(clientResponse -> clientResponse.bodyToMono(LocationNevarisDto.class)) .subscribe(locationDto -> { List<ProjectInfosNevarisDto> projectInfos = new ArrayList<>(locationDto.getProjektInfos()); for (ProjectInfosNevarisDto projectInfosNevarisDto:projectInfos) { this.webClient.get() .uri(projectURI + projectInfosNevarisDto.getId()) .accept(MediaType.APPLICATION_JSON) .exchangeToMono(clientResponse -> clientResponse.bodyToMono(ProjectNevarisDto.class)) .subscribe(projectNevaris -> { if("AF".equals(projectNevaris.getStatus())) { projectService.saveNevarisProject(projectNevaris); } }); } });
涉及DTO定义
@Data public class LocationNevarisDto { private String id; private String bezeichnung; private DatenbankInfo datenbankInfo; @Data public static class DatenbankInfo { private String server; private String datenbank; private String benutzername; } private List<ProjectInfosNevarisDto> projektInfos; }
优化方案代码
// 第一步:获取项目ID列表并启动同步流程 this.webClient.get() .uri(locationURI) .accept(MediaType.APPLICATION_JSON) .retrieve() // 替代exchangeToMono,简化响应处理逻辑 .bodyToMono(LocationNevarisDto.class) // 设置全局超时,避免获取列表阶段卡住 .timeout(Duration.ofSeconds(10)) // 针对超时/网络异常做指数退避重试 .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)) .filter(t -> t instanceof TimeoutException || t instanceof WebClientResponseException)) // 把项目列表转成Flux,方便逐个处理 .flatMapMany(locationDto -> Flux.fromIterable(locationDto.getProjektInfos())) // 控制并发请求数(这里设为5,可根据目标服务能力调整) .flatMap(projectInfo -> fetchAndSaveProject(projectInfo.getId()), 5) // 收集所有处理结果,等待全部操作完成 .collectList() // 全量完成后的通知逻辑 .subscribe( resultList -> { long successCount = resultList.stream().filter(Boolean::booleanValue).count(); System.out.println("所有项目同步完成:成功" + successCount + "个,失败" + (resultList.size() - successCount) + "个"); }, error -> System.err.println("同步流程整体出错:" + error.getMessage()) ); // 封装单个项目的获取+保存逻辑,实现错误隔离 private Mono<Boolean> fetchAndSaveProject(String projectId) { return this.webClient.get() .uri(projectURI + projectId) .accept(MediaType.APPLICATION_JSON) .retrieve() .bodyToMono(ProjectNevarisDto.class) // 单个项目请求的超时设置 .timeout(Duration.ofSeconds(8)) // 单个请求的重试策略 .retryWhen(Retry.backoff(2, Duration.ofMillis(500))) // 过滤仅保存状态为AF的项目 .filter(project -> "AF".equals(project.getStatus())) .flatMap(project -> { projectService.saveNevarisProject(project); return Mono.just(true); }) // 单个项目处理失败时,记录日志且不中断整个流程 .onErrorResume(error -> { System.err.println("处理项目ID[" + projectId + "]失败:" + error.getMessage()); return Mono.just(false); }); }
关键优化点说明
- 并发控制:通过
flatMap的第二个参数限制同时发起的请求数,避免瞬间高并发触发目标服务限流或客户端超时 - 超时与重试:给不同阶段的请求设置独立超时,针对超时、网络异常做指数退避重试,降低失败概率
- 流程编排:用Reactor原生操作符替代嵌套
subscribe,实现异步流程的统一管控 - 错误隔离:单个项目处理失败不会中断整个同步流程,同时记录错误日志便于排查
- 全量完成监听:通过
collectList().subscribe()统一监听所有操作的完成状态,方便后续通知或收尾逻辑
额外建议
- 如果目标服务支持批量查询接口,优先改用批量请求,减少HTTP调用次数
- 若DB保存是同步操作,建议用
Mono.fromCallable()包装成异步操作,避免阻塞Reactor线程 - 可添加Metrics监控请求耗时、成功率,方便定位超时瓶颈
内容的提问来源于stack exchange,提问作者beks6
相关产品推荐
相关产品推荐

