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

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();
}

核心要点说明

  1. 并行性保障:flatMap默认支持异步并行执行,确保所有设备的Set/Get调用同时发起,和你之前用Mono.zip的效果一致,且支持动态新增设备。
  2. 无阻塞逻辑:全程使用Reactor操作符串联流程,没有任何阻塞调用,完全适配响应式线程模型,不会出现block()的报错。
  3. 重试可控:通过repeatWhen精确控制重试触发条件和次数,还能灵活添加重试间隔,避免过度占用设备资源。
  4. 任务完整性:只要这个Mono被正确订阅(比如交给Spring Web处理或后台异步执行),即使前端请求超时,整个循环也会在后台执行完成。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 06:05:53