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保留单节点请求复用能力,在异常场景下兼顾主链路可用性与最低成本控制。