上个月公司做大模型能力聚合网关,业务线产品经理提了一个挺现实的诉求:前端对话框必须做到极致的“打字机秒出”。在他们眼里,用户不管你后面跑的是千亿参数还是量化版本,如果点了发送按钮超过两秒屏幕还没动静,就会被判定为系统卡死。
但当时我们后端接入了四五个不同供应商的 API,既有公有云的主流商用模型,也有本地机房私有化部署的开源大模型。每个厂商的 API 稳定性和推理速度天差地别,同一款模型在早高峰和半夜的响应时间也能差出两三倍。为了做动态路由和智能降级,我们必须在生产环境中持续对各厂商模型进行并行测速。
做大模型测速,最关键的指标不是整句话吐完的总耗时(Total Latency),而是 TTFT(Time to First Token,首字延迟)。只要首个 Token 快速蹦出来,用户的心理等待阈值就会大幅放宽。今天聊聊怎么用 Spring AI 配合 Project Reactor 的Flux,优雅地实现多厂商模型的并行测速与指标聚合。
为什么传统阻塞方式测不准首字延迟?
很多做传统 Java Web 开发的朋友刚转到大模型对接时,习惯性地拿RestTemplate或者普通的 HTTP Client 测接口。发送请求,等整个响应 JSON 包收齐,算个耗时。但大模型推理本质上是自回归生成,首字延迟与后续生成速度是两个完全解耦的维度:
- TTFT 核心影响因素:Prompt 编码(Prefill 阶段)、网络握手耗时、服务方排队时间、KV Cache 命中率。
- 生成速度影响因素:Decode 阶段显卡带宽、Batch Size 大小、生成 Token 总量。
如果用全量响应来评估体验,会出现严重误判。比如 A 模型首字只需 280ms,但生成 1000 字花了 5 秒;B 模型首字卡了 1.8 秒,但后续吐字飞快,2 秒吐完。在客户端交互上,A 模型的实际体验远优于 B 模型。
因此,我们必须使用 SSE(Server-Sent Events)流式响应,并在接收到第一个数据 Chunk 的瞬间打上高精度时间戳。在 Spring 技术栈里,Spring AI 提供的流式客户端天然返回 Project Reactor 的Flux<ChatResponse>,这正是做异步并行统计的利器。
响应式测速流水线设计
我们的目标是:针对同一段测试 Prompt,同时向供应商 A、供应商 B、供应商 C 并行发起调用,利用响应式操作符精准拦截首字时间戳、末字时间戳,以及计算每秒生成 Token 速率(Tokens Per Second, TPS),最终汇总成一份对比数据。
整个流程如果用线程池加CountDownLatch来做,代码会显得非常冗长,并且线程开销和上下文切换也会影响计时精度。而借助Flux的操作符编排,几行声明式代码就能完成并行触发、分流拦截和超时兜底。
来看核心的测速监控结构体定义:
public record ModelSpeedMetric( String providerName, String modelName, long promptTokens, long completionTokens, long dnsAndConnectMs, // 网络建连耗时 long ttftMs, // 首字延迟(从发起请求到第一个有效Token) long totalDurationMs, // 总耗时 double tokensPerSecond// 吐字速率 ) {}接下来是核心测速服务类的实现。我们利用Flux.defer保证每次订阅时生成独立的上下文,通过局部状态追踪首字到达时刻:
@Service public class ModelBenchmarkService { private final Map<String, ChatModel> chatModels; public ModelBenchmarkService(Map<String, ChatModel> chatModels) { this.chatModels = chatModels; } public Mono<ModelSpeedMetric> measureSingleModel(String providerKey, String promptText) { ChatModel chatModel = chatModels.get(providerKey); if (chatModel == null) { return Mono.error(new IllegalArgumentException("未找到对应的模型实例: " + providerKey)); } return Mono.defer(() -> { long requestStartTime = System.nanoTime(); AtomicLong firstTokenTime = new AtomicLong(0); AtomicLong totalTokens = new AtomicLong(0); Prompt prompt = new Prompt(new UserMessage(promptText)); return chatModel.stream(prompt) .doOnNext(response -> { // 记录首字延迟:仅在第一次收到非空文本块时打点 String content = response.getResult().getOutput().getText(); if (content != null && !content.isEmpty()) { firstTokenTime.compareAndSet(0, System.nanoTime()); totalTokens.incrementAndGet(); } }) .doOnError(ex -> { // 记录特定厂商异常日志,便于排查推理引擎报错 }) .then(Mono.fromCallable(() -> { long endTime = System.nanoTime(); long startMs = requestStartTime; long firstMs = firstTokenTime.get(); long ttft = (firstMs > 0) ? (firstMs - startMs) / 1_000_000 : -1; long totalMs = (endTime - startMs) / 1_000_000; long tokens = totalTokens.get(); double tps = (totalMs > ttft && ttft > 0) ? (tokens * 1000.0 / (totalMs - ttft)) : 0.0; return new ModelSpeedMetric( providerKey, chatModel.getClass().getSimpleName(), 0, // 提示词Token预估或从Metadata读取 tokens, 0, ttft, totalMs, Math.round(tps * 100.0) / 100.0 ); })) .timeout(Duration.ofSeconds(15)) .onErrorResume(ex -> Mono.just(new ModelSpeedMetric( providerKey, "TIMEOUT_OR_FAILED", 0, 0, 0, -1, -1, 0.0 ))); }); } public Flux<ModelSpeedMetric> benchmarkAllProviders(List<String> providers, String testPrompt) { // 利用 Flux.merge 并行分发所有厂商的测速任务 return Flux.fromIterable(providers) .flatMap(provider -> measureSingleModel(provider, testPrompt) .subscribeOn(Schedulers.boundedElastic())); } }生产落地的避坑经验
上面的逻辑跑在本地单元测试里看起来非常顺利,但一旦搬到生产环境做定期的自动化探测,有几个非常隐蔽的坑必须提前避开。
1. HTTP 客户端长连接与冷启动握手偏差
测速最大的误差来源往往不是大模型本身,而是 TCP 与 TLS 握手。
如果是第一次请求某个厂商的域名,DNS 解析需要 2050ms,TLS 1.3 握手需要 12 个 RTT(跨洋网络可能达到 200~300ms)。这部分网络耗时会直接混入 TTFT 中,导致探测结果剧烈抖动。
在配置 Spring AI 底层的WebClient时,务必开启 HTTP 连接池保持长连接(Keep-Alive):
ConnectionProvider provider = ConnectionProvider.builder("custom-ai-pool") .maxConnections(50) .maxIdleTime(Duration.ofSeconds(60)) .maxLifeTime(Duration.ofMinutes(5)) .pendingAcquireTimeout(Duration.ofSeconds(5)) .evictInBackground(Duration.ofSeconds(30)) .build(); HttpClient httpClient = HttpClient.create(provider) .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 3000) .responseTimeout(Duration.ofSeconds(30));在正式做对比之前,先发送一个探测心跳(例如 Prompt 为 "hi")做连接预热,把连接建立耗时和推理准备阶段隔离开。
2. 空包与元数据 Chunk 的过滤
不同厂商的流式输出格式规范不完全统一。有的厂商在第一个 SSE 包里只返回role: assistant,content字段是空字符串;有的厂商会在首包里塞入当前的计费元数据。
如果在doOnNext中仅仅判断response != null就打上时间戳,你会发现某些模型的 TTFT 只有惊人的 40ms。但仔细看网络报文,那只是一个空的初始化包,真正的文字内容直到 600ms 后才到达。
因此,代码中的判断条件必须严格校验content != null && !content.trim().isEmpty(),确保打点对应的是人类可见的第一个有效文字。
3. 反应式背压与缓冲导致的滞后
在 Reactive Streams 体系中,下游的处理速度如果慢于上游的推送速度,或者中间经过了某些带有缓冲性质的操作符(如buffer、window),会导致打点时间被人为延后。
在测速链路中,doOnNext应当直接挂在流的最前端(即直接紧随chatModel.stream()之后),切忌在中间插入复杂的日志序列化或数据库写入逻辑。所有重型的持久化操作,都放到流结束之后的then或专门的异步事件队列里异步处理。
最终指标的应用场景
拿到准确的 TTFT 和 TPS 数据后,不要只把它当成监控看板上的几个折线图。我们在网关层基于这些指标落地了两项实用策略:
- 会话级动态路由:对即时性要求高的场景(如客服即时问答、代码自动补全),加权路由算法优先倾斜到最近 5 分钟 TTFT 中位数最低的厂商节点。
- 静默双发(Hedge Request):对重要 VIP 客户的提问,同时向两个不同厂商发起请求。只要任一通道在 600ms 内未产生第一个 Token,立即唤醒备用通道。谁先返回首字就消费谁,并取消另一个未完成的流。
把大模型当成传统的黑盒三方接口来管理是行不通的。利用好响应式编程与精细化时间切片,才能把控住 AI 时代的系统体验底线。