← 返回笔记

learning note

Single-flight 是什么?

记录单 JVM 内短生命周期请求复用器的核心思想:同一个 key 的并发请求只让一个执行,其他等待共享结果。

一句话总结:

Single-flight(单飞模式)是一种并发请求合并技术,它的核心思想是:当多个 goroutine / 线程同时请求同一个资源时,确保只有一个 goroutine / 线程真正执行实际操作,其他所有请求都等待并共享这个结果

我们设计的单体 Single-flight,本质上是一个基于 ConcurrentHashMap + CompletableFutureJVM 内短生命周期请求复用器。它通过共享同一 key 的执行结果,解决了同一实例内 AI 请求并发重复执行的问题,降低了单机内的模型调用成本和结果抖动风险;

但它的能力边界也非常明确——只能在单 JVM 内生效,无法跨节点共享状态、回放结果或进行故障接管,因此最终需要进一步演进为分布式 Single-flight。

什么是 Single-flight

Single-flight 的核心思想很简单:

当多个并发请求同时查询同一个 Key 时,只让第一个请求“飞”出去(Flight)去查数据库或执行下游逻辑,其余的请求在本地阻塞等待;等第一个请求拿到结果后,直接把结果共享给所有等待的请求。

它解决的不是“历史结果缓存”问题,而是“当前正在执行中的重复请求”问题。这点很重要

小牛来给大家举个最简单的例子:

  • 线程 A 调用 AI 评分服务,请求 key 是 k1

  • 在线程 A 还没执行完的时候,线程 B 也发起同一个 k1

  • 如果没有 Single-flight,那么 A 和 B 都会真的去调一次 AI

  • 如果有 Single-flight,那么只有 A 真正调用 AI,B 直接等待 A 的结果

所以 Single-flight 更像是“执行中的请求去重器”,而不是普通缓存。

核心原理

Single-flight 的原理非常简单直观:

  1. 请求分组:使用一个唯一的 key 来标识相同的请求(如缓存 key、接口参数哈希)

  2. 状态跟踪:内部维护一个 map,key 是请求标识,value 是正在执行的请求状态

  3. 结果共享:当新请求到来时,先检查 map 中是否已有相同 key 的请求在执行

    • 如果有:直接等待该请求的结果

    • 如果没有:创建一个新的请求并执行,同时将其加入 map

  4. 结果广播:当请求执行完成后,将结果返回给所有等待的请求,并从 map 中移除该请求

可以用一个生活中的例子来理解:多个人同时想点同一家外卖,与其每个人都打开 APP 下单,不如大家凑在一起,由一个人下单,然后大家共享这份外卖。这就是 Single-flight 的工作方式。

我们的代码:

当前核心实现类在:

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

相关配置在:

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

核心代码:

package com.hewei.hzyjy.xunzhi.interview.application.guard;

import cn.hutool.core.util.StrUtil;
import com.hewei.hzyjy.xunzhi.interview.config.InterviewAiSingleFlightConfiguration;
import io.micrometer.core.instrument.MeterRegistry;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;

import java.util.Objects;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Supplier;

@Service
@RequiredArgsConstructor
@Slf4j
public class InterviewAiSingleFlightService {
// 定义本地 single-flight 服务类。

    private final InterviewAiSingleFlightConfiguration configuration;
    // 读取 single-flight 相关配置,例如是否启用、TTL、多长时间等待等。

    private final MeterRegistry meterRegistry;
    // 指标上报组件,用来记录 hit/miss 指标。

    private final ConcurrentMap<String, FlightEntry> flights = new ConcurrentHashMap<>();
    // 本地内存里的“飞行中请求表”。
    // key:请求指纹
    // value:FlightEntry,里面保存这次请求的 future 和过期时间

    public <T> T execute(String key, Supplier<T> supplier) {
    // 对外暴露的核心方法。
    // key 表示请求唯一标识。
    // supplier 表示真正要执行的逻辑,比如调用 AI。
    // <T> 说明这个方法支持任意类型返回值。

        Objects.requireNonNull(supplier, "supplier cannot be null");
        // 先校验 supplier 不能为空。
        // 因为如果没有 supplier,就没有实际要执行的逻辑。

        if (!Boolean.TRUE.equals(configuration.getEnable()) || StrUtil.isBlank(key)) {
        // 如果配置里没启用 single-flight,或者 key 是空的,
        // 就说明这次请求不能安全参与复用,直接绕过。

            meterRegistry.counter("ai_singleflight_miss_total").increment();
            // 记录一次 miss。
            // 这里 miss 的含义是:没有复用已有 flight,而是自己执行。

            return supplier.get();
            // 直接执行调用方传进来的逻辑,并返回结果。
        }

        long now = System.currentTimeMillis();
        // 记录当前时间,后面会用于判断过期和计算 flight 的过期时间。

        long ttlMillis = resolveTtlMillis();
        // 解析 flight 的 TTL(生存时间)。
        // 表示这个 flight 在本地内存中允许被复用多久。

        AtomicBoolean newFlight = new AtomicBoolean(false);
        // 用一个原子布尔值记录:
        // 当前线程是不是创建了一个全新的 flight。

        FlightEntry entry = flights.compute(key, (ignored, existing) -> {
        // 对 flights 这个并发 map 做原子计算。
        // 好处是:多个线程并发进来时,不会重复创建多个 flight。

            if (existing == null || existing.expireAtMillis <= now) {
            // 如果当前 key 没有对应 flight,
            // 或者虽然有但已经过期了,
            // 那就要创建一个新的 flight。

                newFlight.set(true);
                // 标记当前线程是“新建 flight”的线程,也就是 owner。

                return new FlightEntry(new CompletableFuture<>(), now + ttlMillis);
                // 创建新的 FlightEntry:
                // 1. 一个新的 CompletableFuture,用来承载最终结果
                // 2. 设置过期时间 = 当前时间 + TTL
            }

            return existing;
            // 如果已有 flight 且没过期,就直接复用已有的 entry。
        });

        if (newFlight.get()) {
        // 如果当前线程是新建 flight 的线程,
        // 说明这次真正执行 supplier 的就是它。

            meterRegistry.counter("ai_singleflight_miss_total").increment();
            // 再记录一次 miss。
            // 因为这个线程确实没有复用别人,而是自己执行。

            try {
            // 开始执行真正的业务逻辑。

                T value = supplier.get();
                // 调用方传入的 supplier 真正执行,比如调 AI,并返回结果。

                entry.resultFuture.complete(value);
                // 把执行成功的结果写进 future。
                // 这样其他等待这个 future 的线程都能拿到同样的结果。

                return value;
                // 当前线程自己也返回这个结果。

            } catch (Throwable ex) {
            // 如果执行过程中抛出任何异常。

                entry.resultFuture.completeExceptionally(ex);
                // 把异常也写进 future。
                // 这样等待者不会一直挂住,而是能感知到同样的异常。

                flights.remove(key, entry);
                // 把这个失败的 flight 从 map 中移除。
                // 避免后续请求继续复用一个已经失败的 entry。

                throw ex;
                // 把异常继续向上抛给当前调用者。

            } finally {
            // 不管成功还是失败,最后都执行清理逻辑。

                cleanupExpired(now);
                // 顺手清理掉已经过期的 flight,避免 map 无限增长。
            }
        }

        meterRegistry.counter("ai_singleflight_hit_total").increment();
        // 如果走到这里,说明当前线程不是 owner,
        // 它复用了已有 flight,所以记一次 hit。

        try {
        // 尝试等待 owner 执行完成。

            @SuppressWarnings("unchecked")
            // 因为 resultFuture 里存的是 Object,这里需要强转成 T。
            // 这个注解是为了压制泛型转换警告。

            T reused = (T) entry.resultFuture.get(resolveWaitTimeoutMillis(), TimeUnit.MILLISECONDS);
            // 等待已有 future 完成。
            // 最多等 waitTimeoutMillis 毫秒。
            // 如果 owner 成功,拿到它的结果;
            // 如果 owner 异常,这里会抛 ExecutionException;
            // 如果等待超时,会抛 TimeoutException。

            return reused;
            // 返回复用到的结果。

        } catch (TimeoutException ex) {
        // 如果等待太久还没完成,说明这个共享执行过程已经不可靠了。

            flights.remove(key, entry);
            // 把这个 entry 从 flights 中移除。
            // 避免后面请求还继续挂在这个可能已经卡死的 future 上。

            throw new CompletionException(new RejectedExecutionException("single-flight wait timeout", ex));
            // 抛出包装后的异常,语义是:
            // 这次 single-flight 等待超时,被拒绝继续等待。

        } catch (InterruptedException ex) {
        // 如果等待过程中线程被中断。

            Thread.currentThread().interrupt();
            // 恢复线程中断标记,这是标准写法,避免吞掉中断信号。

            throw new CompletionException(ex);
            // 把异常包装后往上抛。

        } catch (ExecutionException ex) {
        // 如果 owner 执行 supplier 时抛了异常,
        // future.get() 会把它包装成 ExecutionException。

            throw rethrow(ex.getCause());
            // 取出原始 cause,再做统一包装/转换后抛出。
        }
    }

    private RuntimeException rethrow(Throwable cause) {
    // 把 Throwable 统一转成 RuntimeException,方便上层继续抛出。

        if (cause instanceof RuntimeException runtimeException) {
        // 如果本来就是运行时异常。

            return runtimeException;
            // 直接返回,不做额外包装。
        }

        return new CompletionException(cause);
        // 如果不是运行时异常,就包装成 CompletionException 返回。
    }

    private long resolveTtlMillis() {
    // 解析 flight 的 TTL 配置。

        Long configured = configuration.getTtlMillis();
        // 从配置中读取 ttlMillis。

        return configured != null && configured > 0 ? configured : 4000L;
        // 如果配置有效就用配置值,否则默认 4000 毫秒。
    }

    private long resolveWaitTimeoutMillis() {
    // 解析等待已有 future 的超时时间。

        Long configured = configuration.getWaitTimeoutMillis();
        // 从配置中读取 waitTimeoutMillis。

        return configured != null && configured > 0 ? configured : 5000L;
        // 如果配置有效就用配置值,否则默认 5000 毫秒。
    }

    private void cleanupExpired(long nowMillis) {
    // 清理已经过期的 flight,防止内存里的 flights map 越来越大。

        Integer configured = configuration.getCleanupThreshold();
        // 读取清理阈值配置。

        int threshold = configured != null && configured > 0 ? configured : 256;
        // 如果配置有效就用配置值,否则默认阈值 256。

        if (flights.size() < threshold) {
        // 如果当前 flight 数量还没达到阈值,就不做清理。
        // 这是一个性能优化,避免每次执行都遍历 map。

            return;
            // 直接返回。
        }

        flights.entrySet().removeIf(entry -> entry.getValue().expireAtMillis <= nowMillis);
        // 遍历并删除所有已经过期的 entry。
        // 过期判断标准是:expireAtMillis <= 当前时间。
    }

    private record FlightEntry(CompletableFuture<Object> resultFuture, long expireAtMillis) {
    // 定义一个轻量级的内部记录对象。
    // resultFuture:保存这次请求最终的结果或异常
    // expireAtMillis:这次 flight 什么时候过期
    }
}

当前单体版 Single-flight 的实现非常克制,没有引入复杂状态机,也没有引入额外中间件,而是基于 JVM 内存结构完成。

单体 Single-flight 在这个项目里解决了什么

  • 它首先解决的是同机并发重复调 AI。比如同一个评分请求、追问生成请求、简历抽题请求,短时间内有多个线程同时打到同一个实例时,不应该每个线程都去调一次模型。

  • 它的工作方式很直接:内存里维护一个 ConcurrentMap<key, FlightEntry>,FlightEntry 里放一个 CompletableFuture。第一个请求进来时创建 future 并真正执行 supplier;后面的相同 key 请求不再执行,而是等待这个 future 完成,然后直接复用结果或异常。

  • 这对项目里的几个核心 AI 场景都有效:面试评分、追问生成、简历抽题、神态分析。只要这些请求最终落在同一个 JVM,而且 key 一样,就能被收敛成一次真实 AI 调用。

  • 它带来的实际价值有三点:第一,减少同一节点上的重复模型调用;第二,减少线程池和外部 AI 服务的瞬时压力;第三,避免同一份请求在同机并发下返回多份不同 AI 结果,造成解析和流程推进抖动。

设计缺陷:

1 只能在单 JVM 内生效

如果系统前面挂了负载均衡什么是负载均衡,同一个 key 的请求分发到不同实例,每个实例都会各自执行一次,单体 Single-flight 完全看不到彼此的状态。

  • 比如用户一次提交面试答案,请求第一次落到 Node A,因为前端超时又重试一次,第二次落到 Node B。对 Node A 来说这是“第一次见到这个 key”,对 Node B 来说也是“第一次见到这个 key”,于是两边都会真的调 AI。单体 Single-flight 根本拦不住这种跨机重复。

2 结果和状态都只存在内存里

进程重启、节点宕机后,flight 状态全部丢失,其他实例也无法感知。

3 无法做跨节点接管

如果当前执行者在中途挂掉,没有别的实例能够知道它原本在执行什么,也无法安全接手。

4 无法做跨节点结果回放

即使某个实例已经成功执行完,其他实例也拿不到这个内存里的 future 或结果对象。

6 没有业务 stage 级精细策略

它只有 JVM 内的 TTL 和等待超时,没有按 stage 区分 heartbeat、结果 TTL、压缩、L1 缓存等更复杂治理能力。

所以在这个项目中,我们要设计一个分布式Single-Flight来进一步缓解AI接口资源浪费的问题。