ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java AI应用异步化高并发实战:从CompletableFuture到虚拟线程

Java AI应用异步化高并发实战:从CompletableFuture到虚拟线程 做Java后端这些年尤其是最近一年密集折腾AI应用落地之后我越来越觉得把“Java”“AI”“异步化”“高并发”四个词放进同一个技术方案里本身就是一场硬仗。AI调用不再是几十毫秒的数据库查询而是动辄几秒起步、长任务几十秒才返回的重型请求要是继续用传统同步阻塞写法用户一多线程池瞬间被打满服务就像堵死的高架桥谁也别想过。这篇文章想聊的就是怎么把AI应用从“一个请求占一个线程”的泥潭里拉出来用异步化设计把高并发真正扛下来。内容会覆盖线程模型选型、CompletableFuture编排、SSE流式输出、限流熔断、虚拟线程落地以及我从线上故障里整理出的排查经验。适合正在做AI Agent、RAG服务、大模型API封装或者准备把AI能力接入高并发业务系统的后端工程师刚接触异步编程的同学可以从头看老手直接翻第四节捡坑。1. AI 应用为什么绕不开异步化与高并发1.1 AI 推理请求的延迟特征它跟普通接口完全是两回事普通业务接口数据库查询几十毫秒Redis命中几毫秒慢一点的服务间调用也就几百毫秒。但大模型推理服务的延迟完全不在一个量级上。一次完整的大模型请求要经历请求序列化、网络传输、排队等待GPU资源、预填充、逐token自回归解码、响应传输这几个阶段。关键在解码阶段是串行的每生成一个token都需要一次完整的前向传播耗时随文本长度线性增长。所以一个2000字的生成任务跑8到15秒属于正常范围长文档总结任务甚至到分钟级我一点都不惊讶。更麻烦的是AI业务往往不是一次模型调用就完事。以RAG为例用户问一个问题系统先做向量检索、关键词检索再把检索到的文档拼进Prompt最后才调大模型。Agent场景更夸张一轮回答可能包含“规划、调用工具、观察结果、二次推理”中间要调模型三到五次。这意味着一次用户请求的端到端延迟是单次模型耗时的n倍。我自己见过一个Agent项目单轮回答用户要等40多秒链路里串了四次模型调用。这种延迟特征决定了服务端不能按照普通接口的思路来设计你面对的不是一个“稍慢的HTTP请求”而是一群动辄占住连接十几秒的重型任务。1.2 同步阻塞模型线程空等是最大的浪费打个比方就明白了。传统Servlet模型就像一家餐厅里一个服务员从客人点菜开始就站在桌边等着一直等到菜做完端上桌才离开。客人多的时候服务员全被“站等”这件事耗死了后厨再快也翻不了台。大模型调用就是这么个场景Tomcat连接器线程池默认200个线程每个请求占一个线程Controller里同步调HTTP接口去请求模型模型没返回时线程什么都没做就卡在socket读上白白占着资源。算一笔账。假设单次模型调用平均8秒200个线程全部处于等待状态服务端能支撑的最大并发就是200吞吐量大概200/825 QPS。如果你做的是Agent或RAG服务一个用户请求内部还要并发发起多次上游调用瓶颈会更早暴露。线程不只是数量问题一个Java线程默认栈大小约1MB200个线程光栈内存就接近200MB同时还有线程上下文切换调度开销。有人会想“那把线程池调到2000不就行了”2000个线程就是2GB栈内存大量线程同时在等待IO调度和GC都会变得很难看。所以纯靠加线程数解决不了问题正确做法是让“线程数”和“并发任务数”解耦这正是异步化最核心的价值线程只处理本地逻辑等待模型响应的过程交给底层IO事件驱动机制去完成。1.3 哪些场景最吃这套设计结合我接触过的项目四类场景最依赖异步化与高并发设计。第一类是BFF或网关层。很多公司会把大模型能力聚合成统一API出口后端接多家模型供应商网关必须异步才能用有限线程池支撑大量并发AI调用。第二类是Agent编排引擎一个用户请求触发多个步骤步骤之间既有依赖又有可并行分支不异步化整体延迟和资源占用都会失控。第三类是RAG服务向量检索、关键词检索、重排可以并行做最后统一喂给模型这里能省出不少端到端时间。第四类是流式应用例如聊天机器人要维持长连接、持续输出token本质上就是一个高并发的长连接系统。如果你在这类系统里做开发或架构后面的细节基本可以直接对应到你的日常工作。2. 高并发 AI 应用的整体架构设计2.1 Spring Boot 入口该选虚拟线程还是 WebFlux异步化在入口层就决定了成败。当前主流路线我梳理下来有三条可以走。方案编程模型并发能力改造难度适配场景Servlet 3.1异步 DeferredResult同步代码为主局部异步中中存量Spring MVC项目快速改造Spring WebFlux Netty全响应式高高学习曲线陡新项目、纯SSE流式转发、IO密集型Spring Boot 3.2 Java 21虚拟线程同步代码接近异步吞吐高低IO密集型、大量阻塞等待场景最推荐如果你维护的是老项目不能因为加了AI功能就把整个栈推倒重来那就在Tomcat线程池上配合虚拟线程或者局部异步化。如果是新项目而且入口功能就是纯粹的SSE流式转发WebFlux也值得考虑但团队必须接受响应式编程的调试成本。我见过有团队强行上WebFlux结果一个NullPointerException查了三天最后发现是响应式流的下游取消问题。所以这一类选型必须衡量团队消化能力不是越新越好。2.2 统一封装大模型调用第一件事是让接口返回 Future整体架构里一定要把模型调用封装成独立模块。模块内部解决HTTP连接池、超时、重试、序列化对外只暴露同步和异步两套方法同步方法给内部低并发任务用异步方法统一返回CompletableFuture流式方法直接返回响应式流。我一般这样定义抽象层public interface LlmClient { // 同步调用给非高并发内部任务用 ChatResponse chatSync(ChatRequest request); // 异步调用核心业务走这里 CompletableFutureChatResponse chatAsync(ChatRequest request); // 流式输出适配SSE等长连接场景 FluxChatChunk chatStream(ChatRequest request); }HTTP客户端选型上普通Spring项目用OkHttp或Apache HttpClient都行响应式项目用WebClient。这里特别提醒一句很多人异步化没做好是因为忽略了HTTP客户端自己的调度模型。OkHttp有自己的Dispatcher线程池WebClient底层用Netty事件循环连接池的maxIdleConnections、keepAliveTimeout、最大并发请求数都得单独调。如果这些参数不跟着业务并发量走线程池再宽也白搭连接池会先被拖垮。2.3 削峰填谷请求队列、信号量与拒绝策略异步化让线程不空了但上游模型服务商仍然有并发上限。冲破上限的结果不是更快的响应而是429限流或者超时雪崩。所以架构里不能缺“并发闸门”我会把两个手段配合起来用。信号量Semaphore负责限制当前正在执行的模型调用数量。比如你采购的模型网关并发上限是200信号量就设200超出部分先排队排队超过N秒自动失败。有界队列负责削峰请求先进ArrayBlockingQueue队列满时按策略处理。很多生产事故都源于无界队列任务越积越多最后内存先爆掉。更规范的做法是用Resilience4j的Bulkhead舱壁隔离按模型供应商、按请求优先级拆成独立小舱一个供应商变慢不至于让所有流量都陪葬。这块具体参数怎么设我在3.4节详细展开。2.4 缓存让高并发从源头降下来异步化和高并发解决的是“请求打过来怎么扛住”但更高级的做法是让相同请求根本打不到模型上去。模型调用有四个可以缓存的位置完整Prompt结果缓存相同文本、相同参数直接返回上次结果适合知识库问答、FAQ这类高频重复场景。Embedding向量缓存相同文本的embedding结果没必要重复算。前缀缓存部分模型服务商支持prompt前缀缓存当系统提示词和工具定义固定时能显著降低首字延迟。工具/插件结果缓存Agent调用工具返回的结果简单场景可以直接复用。我习惯用Caffeine做本地缓存它对模型结果这种访问延迟敏感的数据特别合适命中一次能省一个网络RTT。不过要注意LLM输出有随机性缓存命中率、key设计、过期策略要根据业务来。比如温度参数大于0的生成结果默认不缓存否则用户可能看到自相矛盾的答案产品上很难解释。3. 核心实现细节用 CompletableFuture 编排 AI 调用3.1 先写一个干净的异步底座CompletableFuture是Java异步编程的主力工具它比我前几年用的Futureget()强太多支持编排、组合、超时、异常处理。最简单的异步调用长这样Service public class LlmService { private final ExecutorService llmExecutor Executors.newFixedThreadPool(64); private final LlmClient llmClient; public CompletableFutureString generateAsync(String prompt) { return CompletableFuture.supplyAsync(() - llmClient.chatSync(buildRequest(prompt)), llmExecutor) .orTimeout(15, TimeUnit.SECONDS) .exceptionally(ex - handleError(prompt, ex)); } }这里三个关键点。supplyAsync指定自定义线程池不要用公共的ForkJoinPool业务线程池隔离能避免互相干扰。orTimeout(15, SECONDS)给异步任务加了超时兜底这是异步编程最容易漏的一环。exceptionally做降级把异常转换成用户能接受的回退文案而不是让异常直接炸到上层。这个底座是所有AI调用编排的基础后面所有复杂逻辑都从它长出来。3.2 RAG 场景把并行度打满RAG流程是最典型的能通过编排吃透并行的场景。一个用户问题进来向量检索和关键词检索互不依赖完全可以并行。串行写的话总耗时是两组检索时间相加再加模型生成时间并行写的话总耗时只等于较慢那组检索时间加模型生成时间。public CompletableFutureString ragAnswer(String question) { CompletableFutureListDocument vectorFuture CompletableFuture.supplyAsync(() - vectorSearch(question), searchExecutor); CompletableFutureListDocument keywordFuture CompletableFuture.supplyAsync(() - keywordSearch(question), searchExecutor); return vectorFuture .thenCombine(keywordFuture, (vec, keyword) - mergeDocuments(vec, keyword)) .thenApplyAsync(docs - llmClient.chatSync(buildPrompt(question, docs)), llmExecutor) .orTimeout(20, TimeUnit.SECONDS) .exceptionally(ex - fallbackAnswer(question)); }thenCombine把两个检索结果合并thenApplyAsync在拿到文档后拼Prompt再去调模型。注意我每一步都指定了线程池检索任务走检索线程池模型调用走模型线程池。这样做的好处是不同任务的线程池可以独立扩容、独立监控模型变慢不会拖住检索线程。等两个调用都完成用户感受到的延迟就是“最长那条链路的耗时”而不是所有任务耗时的加总。3.3 流式 SSE 输出怎么做聊天类AI应用几乎都要求流式输出效果上用户能看着字蹦出来心理等待时间从10秒降到1秒。服务端实现我推荐直接上WebFluxController返回FluxSseEvent配合上游模型的SSE流逐token转发。GetMapping(value /chat, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxString chat(RequestParam String prompt) { return Flux.concat( Flux.just(开始处理\n), Flux.from(streamCompletion(prompt)) .map(Chunk::getText) .doOnNext(text - log.debug(token: {}, text)) ) .timeout(Duration.ofSeconds(60)) .onErrorResume(ex - Flux.just(出错了 ex.getMessage())); } private FluxChunk streamCompletion(String prompt) { return webClient.post() .uri(/v1/chat/completions) .bodyValue(streamRequestBody(prompt)) .retrieve() .bodyToFlux(Chunk.class); }流式场景有三个坑必须提前排掉。一是反压上游模型吐token速度可能比客户端消费快如果不做缓冲控制内存会被未消费的token堆满。二是心跳长连接超过几十秒没数据中间网关或客户端可能认为连接断了需要定时发心跳数据保活。三是客户端断连检测用户中途关掉页面服务端必须及时取消上游订阅否则模型还在消耗算力和连接白花钱还占资源。3.4 超时、限流、熔断一个不能少异步化之后系统最怕的不是慢而是“慢却不报错”。我自己的量化方法很简单给每一层超时都取包络值让整体超时时间严格小于上游能接受的极限。超时层次建议值说明HTTP Client读超时15秒超过就中断连接避免线程无限挂起CompletableFuture.orTimeout20秒业务整体兜底含排队时间流式总超时60秒长文本生成也要有上限队列等待超时5秒排队太久的请求直接返回繁忙提示限流上本地用Semaphore控制并发分布式场景上Sentinel或Resilience4j。熔断要关注三个参数失败阈值比如10秒内失败率超过50%就熔断、熔断时长建议30到60秒、半开探测请求数先放1个请求试探上游恢复没有。重试策略只在幂等场景开并且用指数退避加随机抖动避免重试风暴把所有流量又打回到刚要恢复的上游。3.5 Java 21 虚拟线程我用还是不用如果项目已经升级到Java 21虚拟线程几乎是AI应用异步化的最优解。它最大的特点是“同步代码、异步效果”你仍然写阻塞风格代码但每个请求独占一个虚拟线程虚拟线程的创建和切换开销极低底层由JVM自动悬挂和恢复。之前200个平台线程撑不住的局面换成一万个虚拟线程都很轻松。ExecutorService executor Executors.newVirtualThreadPerTaskExecutor(); public CompletableFutureChatResponse generateAsync(String prompt) { return CompletableFuture.supplyAsync(() - llmClient.chatSync(buildRequest(prompt)), executor) .orTimeout(15, TimeUnit.SECONDS); }但要冷静看两件事。第一虚拟线程并不意味着可以无限发请求上游模型并发上限还在必须用Semaphore闸门把同时调用的虚拟线程数量限制住否则上游429会教你做人。第二Java 21早期版本里虚拟线程如果在synchronized块或高频ThreadLocal操作中被钉扎会占用底层平台线程影响并发效果。规避方案是能不用synchronized就不用换成ReentrantLockThreadLocal改用TransmittableThreadLocal并注意清理。我的建议是存量系统优先虚拟线程新项目如果团队对响应式没有把握也优先虚拟线程WebFlux只在纯粹的流式网关场景再用。4. 常见问题与排查技巧实录4.1 流量一高线程池就满问题到底在哪现象很好认日志里大量RejectedExecutionException线程池监控显示活跃线程数长期打满。很多人第一反应是“线程池开小了”于是把线程数翻倍结果过两天又满。这时候真正要查的是“线程到底在等什么”。用jstack抓一份线程栈如果大量线程停在socketRead0这种网络读等待上说明瓶颈在上游模型接口慢或者连接池借不到连接如果大量线程停在队列取任务说明本地逻辑处理能力确实不够。我的排查顺序是先看线程状态分布再看HTTP连接池指标最后看模型上游的响应P99。经验是AI场景下线程池被打满80%的原因是上游模型慢导致线程全在等IO20%才是线程配置问题。优化方向通常是限制上游并发而不是无限加线程。4.2 TraceId在异步线程里断了排查变成灾难这是异步化最烦人的问题。同步代码里用MDC.put(traceId, ...)日志自然串起来一上异步子任务跑到别的线程ThreadLocal里的TraceId全丢了。线上出了问题几十个节点的日志对不上号等于瞎查。解决方案有两个。一是用阿里的TransmittableThreadLocal配合TtlExecutors.getTtlExecutorService包装线程池值会自动在父线程和子线程间传递。二是如果不想引依赖就手动在提交任务时把MDC上下文快照传进去任务执行前再放回去。public T CompletableFutureT submitWithMdc(CallableT task, ExecutorService executor) { MapString, String context MDC.getCopyOfContextMap(); return CompletableFuture.supplyAsync(() - { if (context ! null) { MDC.setContextMap(context); } try { return task.call(); } finally { MDC.clear(); } }, executor); }我建议无论用哪种方案在写第一行异步代码时就把它加上否则等线上出问题再补日志已经丢了无数条。4.3 HTTP连接池和超时叠加慢查询拖垮全链路还有一类故障表现是线程没满但大量请求卡在“获取连接”阶段。比如WebClient默认连接池是固定大小某个慢请求把连接借走之后一直不还后来的请求全部排队等连接。这时候看连接池指标池里空闲连接数接近0等待获取连接线程数很高。根子在于两层超时叠加出了问题HTTP读超时设得比业务整体超时长业务超时到了之后连接拿不到读超时依然没到于是资源被无效请求占着。我的调法是把“读超时”设为“业务超时”的三分之二让底层比上层更早放弃同时限制单路由连接池大小并设空闲回收。连接池大小要和并发信号量做联动信号量限制了整体并发连接池大小就按上限加一点缓冲别让两者互相独立、互相拆台。4.4 SSE流式连接被卡死网关和客户端都有份流式场景线上最常见的现象是客户端等半天第一个字都没出来或者输出到一半突然不动了。先查网关Nginx默认proxy_read_timeout是60秒如果模型生成慢60秒内没有任何新数据网关直接掐断。这个值要调大并结合应用层心跳。再查服务端确认有没有主动定期发送心跳数据最后查客户端前端如果用了代理或切换过网络可能连接早就断了但服务端不知道不会自动取消生产任务。我踩过一次很深的坑SSE请求量一大服务端每个流式请求都占着一个连接和一个虚拟线程但客户端突然大批量关页面服务端却没及时收手模型调用还在继续。后来加了onCancel回调客户端断连就立刻取消订阅再配合总活跃流数量上限问题才算根治。4.5 内存、线程与GC监控这几条就能活命异步化和高并发改造之后系统资源和原来完全不是一个画像。原来线程数稳定内存消耗稳定现在是线程数可能波动极大、任务队列忽高忽低、HTTP连接池频繁借还。我建议至少盯住四个指标活跃线程数、任务队列深度、HTTP客户端连接池活跃连接数、JVM堆内存和GC耗时。任何一项连续五分钟出现异常增长都要当事故处理。GC方面有一个观察点如果用的是虚拟线程ThreadLocal里堆积的数据会成为内存泄漏的老鼠仓因为虚拟线程数量极大每个里面都挂一份上下文内存膨胀会比平台线程更隐蔽。线上一定要定期清理ThreadLocal并对MDC、TTL封装做完整性测试。说到底高并发设计不只是把代码写成异步就完事资源监控和问题预案得一起配套上线才行。最后分享一个我自己的土办法所有异步任务的日志里一律把超时时间、当前并发水位、队列长度打进去。线上排查时这三条数字能直接定位是超时配置问题、上游变慢问题还是排队问题比翻半天监控面板快得多。这套异步化方案从设计到落地我自己前后迭代了四五轮每次线上故障都在提醒我“等待不可怕失控才可怕”。如果你的AI应用也正在被并发压得喘不过气按这条链路走下去方向不会错。
返回列表