Spring Boot响应式API实现下游设备set-get重试逻辑的技术问题
解决方案:纯响应式实现Set->Get循环重试
绝对不要在响应式线程中使用block()或CountDownLatch这类阻塞操作,Reactor提供了完整的操作符来实现等待、判断和重试逻辑,完全贴合响应式编程模型。
步骤1:封装设备API调用(保留并行性)
先把单设备的Set/Get调用封装成Mono,再通过Flux.flatMap实现多设备并行调用:
// 单设备Set调用,返回是否设置成功 private Mono<Boolean> setDevice(String deviceName, int targetValue) { String url = String.format("http://%s/devicename/set/%d", deviceName, targetValue); return webClient.get() .uri(url) .retrieve() .bodyToMono(String.class) .map(resp -> resp.contains("成功")); // 根据实际返回结果判断成功与否 } // 单设备Get调用,返回当前值 private Mono<Integer> getDeviceValue(String deviceName) { String url = String.format("http://%s/devicename/get", deviceName); return webClient.get() .uri(url) .retrieve() .bodyToMono(Integer.class); } // 并行调用所有设备的Set,返回所有设置结果 private Mono<List<Boolean>> batchSetAllDevices(List<String> deviceNames, int targetValue) { return Flux.fromIterable(deviceNames) .flatMap(device -> setDevice(device, targetValue)) .collectList(); } // 并行调用所有设备的Get,返回所有当前值 private Mono<List<Integer>> batchGetAllDevices(List<String> deviceNames) { return Flux.fromIterable(deviceNames) .flatMap(this::getDeviceValue) .collectList(); }
步骤2:实现单次Set->Get流程
用then操作符等待所有Set完成后执行Get,再判断设备值是否合规:
// 判断所有设备值是否符合目标要求 private boolean allDevicesCompliant(List<Integer> currentValues, int targetValue) { return currentValues.stream().allMatch(v -> v == targetValue); } // 单次Set->Get完整流程,返回是否合规 private Mono<Boolean> singleSetGetCycle(List<String> deviceNames, int targetValue) { return batchSetAllDevices(deviceNames, targetValue) // 等待所有Set调用完成后,执行批量Get .then(batchGetAllDevices(deviceNames)) // 判断结果是否合规 .map(values -> allDevicesCompliant(values, targetValue)); }
步骤3:实现循环重试逻辑
用repeatWhen操作符控制重试次数,直到所有设备合规或重试耗尽:
public Mono<Void> runSetGetLoop(List<String> deviceNames, int targetValue, int maxRetries) { AtomicInteger retryCounter = new AtomicInteger(0); return singleSetGetCycle(deviceNames, targetValue) // 仅当不合规时触发重试 .repeatWhen(completed -> completed .flatMap(isCompliant -> { if (isCompliant) { // 合规,终止循环 return Mono.empty(); } int currentRetry = retryCounter.incrementAndGet(); if (currentRetry >= maxRetries) { // 重试次数耗尽,终止循环 return Mono.empty(); } // 可选:添加重试延迟,避免频繁调用设备API return Mono.delay(Duration.ofSeconds(5)); })) // 忽略最终结果,返回Void表示任务完成 .then(); }
核心要点说明
- 并行性保障:
flatMap默认支持异步并行执行,确保所有设备的Set/Get调用同时发起,和你之前用Mono.zip的效果一致,且支持动态新增设备。 - 无阻塞逻辑:全程使用Reactor操作符串联流程,没有任何阻塞调用,完全适配响应式线程模型,不会出现
block()的报错。 - 重试可控:通过
repeatWhen精确控制重试触发条件和次数,还能灵活添加重试间隔,避免过度占用设备资源。 - 任务完整性:只要这个
Mono被正确订阅(比如交给Spring Web处理或后台异步执行),即使前端请求超时,整个循环也会在后台执行完成。
内容的提问来源于stack exchange,提问作者Brad Nelson
相关产品推荐
相关产品推荐

