← 返回笔记

learning note

本地降级机制

记录分布式协调不可用或 Redis 异常时,如何用本地降级维持核心请求链路的可用性。

背景:

这里主要就做一个容灾机制而已。我们的分布式Single-flight依赖的是Redis来搭建状态机。但一个成熟的分布式方案,不能只考虑“理想情况下如何工作”,还必须考虑对应的容灾情况。

这就是我们本地降级机制存在的必要性。也就是说,当 Redis 协调、Lua 状态机、分布式回放或接管逻辑不能稳定工作时,系统并不会直接全线失败,而是可以:

  • 退回到本机内的 single-flight

其实说人话就是:当一个请求打到服务上后,我们会优先使用分布式Single-Flight。而如果分布式Single-Flight崩盘的时候,我们会回退到普通的单机Single-Flight上。

** 所以这个降级机制没什么高大上的,他就是一个单机版本的Single-Flight而已。**

单机** single-flight**** 解决的是什么**

单机版本的 single-flight,核心能力是:

  • 在同一个 JVM 内,如果同一个 key 的请求并发到来,只让一个 leader 真正执行,其他线程复用同一个 CompletableFuture

当前项目里的单机实现类是:

  • admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/singleflight/service/InterviewAiSingleFlightService.java

它本质上依赖的是:

  • ConcurrentHashMap

  • compute(...)

  • CompletableFuture

所以它能解决的是:同一个服务器节点内的请求风暴。

当前代码里的本地降级是怎么设计的

相关核心代码在:

  • admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/singleflight/service/DistributedInterviewAiSingleFlightService.java

先看入口:

public String execute(String stage, String requestKey, Supplier<String> supplier)

这个方法里,本地降级大致有三层。

第一层:本来就不走分布式,直接走本地

代码里一开始就有:

if (!Boolean.TRUE.equals(configuration.getEnable())
        || mode == FlightMode.LOCAL
        || !Boolean.TRUE.equals(configuration.getDistributedEnabled())) {
    return localSingleFlightService.execute(requestKey, supplier);
}

这表示,在下面这些情况下,系统压根不尝试分布式协调:

  • 总开关没开

  • 运行模式就是 LOCAL

  • 分布式能力未启用

此时执行的不是“失败后降级”,而是:策略层面明确只用本地 single-flight

**第二层:分布式执行抛异常,在 ****HYBRID**模式下回退本地

代码继续往下看:

try {
    return executeDistributed(stage, requestKey, supplier);
} catch (RuntimeException ex) {
    if (mode == FlightMode.HYBRID) {
        log.warn("Distributed single-flight fallback to local mode ...");
        return localSingleFlightService.execute(requestKey, supplier);
    }
    throw ex;
}

这才是狭义上最典型的:本地降级。它的语义是:

  • 我优先尝试分布式 single-flight

  • 但如果分布式链路抛了运行时异常

  • 且当前模式是 HYBRID

  • 那我就退回本地 single-flight

第三层:分布式内部遇到异常结果,直接走本地

除了最外层 try-catch,在 executeDistributed(...) 里还有两处很有意思的本地回退。

第一处:

if (acquireResult == null || acquireResult.getAction() == null) {
    return localSingleFlightService.execute(safeRequestKey, supplier);
}

这表示:

  • 如果 Lua 协调结果拿不到

  • 或者协议解析后没有合法 action

  • 那就不要继续在分布式状态机里硬跑

  • 直接退回本地执行

第二处:

default -> {
    return localSingleFlightService.execute(safeRequestKey, supplier);
}

这表示:

  • 即使拿到了 FlightAcquireResult

  • 但 action 落入了未知分支

  • 也不继续冒险,而是直接走本地兜底

这个设计很成熟,因为它体现了一个原则:对分布式状态机的不确定性,宁可保守回退,也不要继续猜测执行

单机Single-Flight设计思路:

代码位置和整体结构

核心实现类:

  • admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/singleflight/service/InterviewAiSingleFlightService.java

相关配置类:

  • admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/config/InterviewAiSingleFlightConfiguration.java

这套实现的骨架非常清晰:

  • 一个 ConcurrentMap<String, FlightEntry> 保存当前 JVM 内正在执行的 flight

  • 一个统一入口 execute(String key, Supplier<T> supplier)

  • 同 key 下通过 compute(...) 原子选出 leader

  • leader 真正执行 supplier,并把结果广播给 waiter

  • waiter 通过等待同一个 CompletableFuture 来复用结果

如果只用一句话概括它的结构,我会这样说:

用一个按 key 管理的 in-flight map,把“相同请求并发执行”收敛成“一个 leader 执行 + 多个 waiter 等待同一个 future”。

这个真没啥难的,在这里就不水字数了。看不懂代码的同学自己问一下豆包

public <T> T execute(String key, Supplier<T> supplier) {
        Objects.requireNonNull(supplier, "supplier cannot be null");
        if (!Boolean.TRUE.equals(configuration.getEnable()) || StrUtil.isBlank(key)) {
            meterRegistry.counter("ai_singleflight_miss_total").increment();
            return supplier.get();
        }

        long now = System.currentTimeMillis();
        long ttlMillis = resolveTtlMillis();
        AtomicBoolean newFlight = new AtomicBoolean(false);
        // compute 保证同 key 下“创建 flight + 复用 flight”原子化,避免瞬时并发下出现多个 leader。
        FlightEntry entry = flights.compute(key, (ignored, existing) -> {
            if (existing == null || existing.expireAtMillis <= now) {
                newFlight.set(true);
                return new FlightEntry(new CompletableFuture<>(), now + ttlMillis);
            }
            return existing;
        });

        if (newFlight.get()) {
            meterRegistry.counter("ai_singleflight_miss_total").increment();
            try {
                // leader 执行真实调用,并把结果广播给同 key 的等待者。
                T value = supplier.get();
                entry.resultFuture.complete(value);
                return value;
            } catch (Throwable ex) {
                entry.resultFuture.completeExceptionally(ex);
                flights.remove(key, entry);
                throw ex;
            } finally {
                cleanupExpired(now);
            }
        }

        meterRegistry.counter("ai_singleflight_hit_total").increment();
        try {
            @SuppressWarnings("unchecked")
            T reused = (T) entry.resultFuture.get(resolveWaitTimeoutMillis(), TimeUnit.MILLISECONDS);
            return reused;
        } catch (TimeoutException ex) {
            // waiter 超时后主动剔除旧 flight,避免后续请求持续等待一个可能已失活的 future。
            flights.remove(key, entry);
            throw new CompletionException(new RejectedExecutionException("single-flight wait timeout", ex));
        } catch (InterruptedException ex) {
            Thread.currentThread().interrupt();
            throw new CompletionException(ex);
        } catch (ExecutionException ex) {
            throw rethrow(ex.getCause());
        }
    }

面试里怎么讲这个点

这个点很适合讲,因为它能体现你不是只会堆 Redis 和 Lua,而是真的考虑过:

  • 分布式方案出问题时系统怎么办

如果你口头讲给面试官,我建议这样表达:

在分布式 AI single-flight 方案里,我没有把主业务完全绑死在 Redis 协调层上,而是设计了本地降级机制。系统优先走 Redis Lua + 状态机做跨节点去重,但在 HYBRID 模式下,如果分布式协调链路抛异常、返回非法动作或无法稳定裁决,就自动回退到本机 single-flight,继续用 ConcurrentHashMap + CompletableFuture 做单节点内的请求合并。这样虽然暂时放弃了跨节点强复用,但能保住单机去重和主链路可用性,避免系统在协调层异常时直接裸奔或整体失败。

如果要写成更简历化的一句,可以写成:

  • 设计分布式 single-flight 本地降级机制,在 HYBRID 模式下于 Redis/Lua 协调异常时自动回退本机 single-flight,基于 ConcurrentHashMap + CompletableFuture 保留单节点请求复用能力,在异常场景下兼顾主链路可用性与最低成本控制。