尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Java 21虚拟线程在RAG平台中的实践:从响应式到同步的性能优化

Java 21虚拟线程在RAG平台中的实践:从响应式到同步的性能优化 1. 项目缘起当传统RAG遇上Java 21的虚拟线程去年年底我们团队负责维护的一个内部RAG检索增强生成平台开始频繁告警。这个平台主要服务于产品文档和客服知识库的智能问答高峰期并发查询能达到每秒上百次。原有的技术栈是基于Spring Boot WebFlux的响应式编程模型配合一个线程池来处理向量检索、LLM调用等IO密集型任务。理论上响应式模型应对高并发是没问题的但实际运维中我们遇到了几个头疼的问题一是响应式编程的学习曲线陡峭团队里能熟练调试复杂异步链的同事不多出了问题排查周期长二是在处理一些需要顺序执行、且有状态依赖的复杂检索逻辑时响应式代码写起来异常别扭可读性急剧下降三是线程池的配置成了玄学设大了浪费资源设小了在流量尖峰时任务排队导致整体响应延迟飙升。就在我们为线程池参数和回调地狱头疼时Java 21正式发布了其核心特性之一——虚拟线程Virtual Threads引起了我的注意。官方宣称它能以极低的开销支持海量并发编写方式还是我们熟悉的同步阻塞式代码。这听起来像是一剂对症良药能否用虚拟线程重写这个RAG平台用同步的写法获得异步的性能同时提升代码的可维护性这个想法让我很兴奋但也知道从架构设计到落地肯定布满荆棘。这篇文章我就来完整复盘这次重写之旅从顶层架构的重新思考到具体编码中的“踩坑”与“填坑”希望能给面临类似技术选型困境的伙伴们一个真实的参考。2. 架构重塑面向虚拟线程的RAG平台设计重写不是简单的替换首先要回答的是在虚拟线程的新范式下整个RAG平台的架构应该如何调整才能最大化其优势同时规避潜在风险2.1 核心组件与数据流再梳理我们的RAG流程是标准的多阶段流水线用户查询 - 查询理解/改写 - 向量库检索 - 多路召回结果融合 - 重排序 - 上下文构建 - 大模型生成 - 响应返回。在旧架构中每个阶段都可能涉及远程调用数据库、向量库、大模型API这些IO操作通过响应式操作符异步组合。新架构的核心转变在于我们将每个处理阶段封装为一个独立的、可组合的“任务单元”这些任务单元不再返回Mono或Flux而是直接返回业务对象。它们由虚拟线程来执行。这样一来整个处理链路在代码层面就是一连串同步方法调用逻辑清晰直白。例如一个简化的核心服务类看起来会是这样的Service public class RagQueryService { private final QueryUnderstandingService understandingService; private final VectorRetrievalService retrievalService; private final RerankService rerankService; private final LlmGenerationService generationService; public RagResponse handleQuery(String userQuery) { // 1. 查询理解 (可能调用NLP服务) EnhancedQuery enhancedQuery understandingService.understand(userQuery); // 2. 向量检索 (IO操作访问向量数据库如Milvus/Pinecone) ListRetrievedChunk chunks retrievalService.retrieve(enhancedQuery); // 3. 重排序 (可能调用重排序模型API) ListRerankedChunk reranked rerankService.rerank(chunks, enhancedQuery); // 4. 上下文构建与LLM生成 (调用OpenAI、DeepSeek等API) String answer generationService.generateAnswer(enhancedQuery, reranked); return new RagResponse(answer, reranked); } }关键设计点handleQuery方法本身是同步的但其中调用的每一个service方法内部凡是涉及阻塞IO的操作如HTTP客户端调用、数据库JDBC操作我们都确保它们运行在虚拟线程上。如何确保这依赖于我们选用的客户端库是否支持。对于不支持的部分我们需要进行封装。2.2 线程模型与并发控制策略这是重写的核心。虚拟线程是“廉价”的我们可以创建成千上万个而无需担心传统操作系统线程的资源消耗。但这并不意味着可以无节制地创建。我们的策略是摒弃复杂的业务线程池拥抱结构化并发Structured Concurrency。Java 21的java.util.concurrent包为结构化并发提供了StructuredTaskScope。这允许我们将一组相关的虚拟线程任务作为一个整体来管理具备以下优势错误传播与取消如果主任务或任何一个子任务失败所有其他子任务会被自动取消避免资源泄漏。清晰的代码结构任务之间的层级和依赖关系在代码中一目了然。例如在“多路召回”阶段我们可能同时发起基于向量的语义检索、基于关键词的全文检索甚至调用外部知识图谱API。使用StructuredTaskScope可以这样实现public ListRetrievedChunk multiPathRetrieve(EnhancedQuery query) { try (var scope new StructuredTaskScope.ShutdownOnFailure()) { // 提交多个检索子任务 SupplierListRetrievedChunk vectorTask scope.fork(() - vectorRetriever.retrieve(query)); SupplierListRetrievedChunk keywordTask scope.fork(() - keywordRetriever.retrieve(query)); SupplierListRetrievedChunk graphTask scope.fork(() - knowledgeGraphRetriever.retrieve(query)); scope.join(); // 等待所有子任务完成 scope.throwIfFailed(); // 如果有任何失败抛出异常 // 合并结果 return fusionStrategy.fuse(vectorTask.get(), keywordTask.get(), graphTask.get()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException(Retrieval interrupted, e); } }这里的一个深刻教训是虽然虚拟线程便宜但StructuredTaskScope内创建的每个fork仍然是一个独立的线程。如果“多路召回”的路由非常多比如超过1000个虽然虚拟线程能创建但下游服务如向量数据库可能无法承受如此突发的并发连接。因此我们引入了每路由的并发限制器使用Semaphore或RateLimiter来控制对同一下游服务的最大并发请求数。2.3 与现有技术栈的融合Spring Boot 3.2我们项目基于Spring Boot。幸运的是Spring Boot 3.2对虚拟线程提供了很好的支持。启用虚拟线程在application.properties中设置spring.threads.virtual.enabledtrueSpring MVC和WebFlux的底层就会使用虚拟线程作为执行器。数据库连接池这是关键。传统的连接池如HikariCP是为平台线程设计的一个物理连接被一个平台线程占用。在虚拟线程场景下如果虚拟线程在等待数据库响应时被挂起它底层的平台线程可以被其他虚拟线程使用但那个数据库连接仍然被“占用”着。如果并发虚拟线程数远大于连接池大小就会导致大量虚拟线程等待连接形成逻辑上的“连接饥饿”。我们的解决方案是适当调大数据库连接池的最大尺寸例如设置为(最大并发请求数 * 每个请求平均持有连接时间)。同时必须确保所有JDBC操作都是真正的阻塞IO并且驱动支持如PostgreSQL的pgjdbc驱动从42.7.0开始支持。HTTP客户端我们使用WebClientSpring的响应式客户端和新的同步RestClient。对于虚拟线程更推荐使用RestClient进行同步调用因为它能自然地让虚拟线程在IO时挂起。如果必须使用WebClient则需要通过block()方法将其响应式结果同步化但这需要小心处理避免在事件循环线程上调用block()。3. 核心实现关键服务与虚拟线程的适配架构确定后进入具体的服务实现。这里分享三个核心服务的适配细节和遇到的坑。3.1 向量检索服务连接池与超时控制我们使用Milvus作为向量数据库。官方提供了Java SDK但其底层HTTP客户端通常是OkHttp的连接池管理需要特别关注。问题初期我们直接使用SDK的默认配置在压力测试时当虚拟线程并发数达到500左右出现了大量的TimeoutException和ConnectionPoolTimeoutException。原因在于OkHttp的默认连接池较小最大5个空闲连接而虚拟线程并发高瞬间创建大量请求虽然线程可以挂起但HTTP连接需要排队等待导致超时。解决方案显式配置OkHttpClient创建自定义的OkHttpClient实例传递给Milvus SDK。OkHttpClient okHttpClient new OkHttpClient.Builder() .connectionPool(new ConnectionPool(200, 5, TimeUnit.MINUTES)) // 增大连接池 .connectTimeout(Duration.ofSeconds(10)) .readTimeout(Duration.ofSeconds(30)) // 向量检索可能较慢 .writeTimeout(Duration.ofSeconds(10)) .build(); // 用这个client初始化MilvusClient在服务层添加熔断与降级使用Resilience4j为检索服务添加熔断器Circuit Breaker。当失败率达到阈值时快速失败并可以降级为返回缓存结果或更简单的关键词检索结果避免雪崩。虚拟线程内的阻塞检测使用Java 21的Thread.currentThread().isVirtual()来确认代码运行在虚拟线程上。我们在关键IO操作前后加了日志确认虚拟线程确实在IO时被挂起而不是意外地在平台线程上阻塞。3.2 LLM生成服务应对长耗时与流式响应调用大模型API如OpenAI GPT-4、DeepSeek是另一个重IO操作而且耗时可能长达数十秒。此外为了用户体验我们常常希望支持流式响应Streaming让答案一个字一个字地返回。挑战如何在同步的虚拟线程模型中处理流式响应解决方案我们采用了“生产者-消费者”模型结合BlockingQueue和虚拟线程。启动一个虚拟线程专门负责调用LLM API并开启流式接收。这个线程将收到的每个数据块chunk放入一个BlockingQueue中。主线程另一个虚拟线程则从BlockingQueue中依次取出数据块并可以通过SSEServer-Sent Events或WebSocket实时推送给前端。public StreamString streamGenerateAnswer(String prompt, ListChunk contexts) { BlockingQueueString queue new LinkedBlockingQueue(); AtomicReferenceException error new AtomicReference(); // 启动生产者虚拟线程 Thread.ofVirtual().start(() - { try { LlmStreamClient client new LlmStreamClient(); client.streamCompletion(prompt, contexts, chunk - { queue.put(chunk); // 将块放入队列 }); queue.put([DONE]); // 流结束标志 } catch (Exception e) { error.set(e); queue.put([ERROR]); } }); // 返回一个从队列中拉取的流这里简化表示 return Stream.generate(() - { try { String item queue.take(); if ([ERROR].equals(item)) throw new RuntimeException(Stream error, error.get()); if ([DONE].equals(item)) return null; // 流结束 return item; } catch (InterruptedException e) { Thread.currentThread().interrupt(); return null; } }).takeWhile(Objects::nonNull); }注意这里queue.take()是阻塞操作它会挂起消费者虚拟线程直到队列中有数据。这正是虚拟线程发挥优势的地方——挂起代价极小。3.3 重排序服务CPU密集型任务的隔离RAG的重排序阶段有时会使用轻量级的交叉编码器模型如bge-reranker在本地进行推理这是一个CPU密集型计算。虚拟线程的陷阱虚拟线程在遇到阻塞IO如socket.read时会自动让出载体线程平台线程但在执行CPU密集型运算时它会一直占用载体线程就像平台线程一样。如果大量虚拟线程同时执行重排序计算会快速耗尽载体线程池默认大小为CPU核心数导致所有虚拟线程包括那些正在等待IO的都被阻塞系统吞吐量不升反降。解决方案将CPU密集型任务与IO密集型任务隔离执行。 我们创建了一个独立的、基于固定大小线程池的ExecutorService专门用于处理重排序计算。Component public class RerankService { // 一个固定大小的线程池大小与CPU核心数相关 private final ExecutorService cpuBoundExecutor Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors()); public ListRerankedChunk rerank(ListRetrievedChunk chunks, EnhancedQuery query) { // 将计算任务提交到独立的线程池 FutureListRerankedChunk future cpuBoundExecutor.submit(() - doRerankComputation(chunks, query)); try { return future.get(); // 这里会阻塞虚拟线程但载体线程被释放 } catch (InterruptedException | ExecutionException e) { throw new RuntimeException(Rerank failed, e); } } private ListRerankedChunk doRerankComputation(ListRetrievedChunk chunks, EnhancedQuery query) { // ... 密集的模型推理计算 ... } }这样负责调度的虚拟线程在future.get()处被挂起载体线程得以释放去服务其他虚拟线程。而实际的计算工作则由专门的平台线程池承担互不干扰。这是混合使用虚拟线程和平台线程的典型场景。4. 性能调优与深度踩坑实录系统跑起来只是第一步性能达标和稳定才是终极目标。这一部分是我们投入精力最多、踩坑最密集的地方。4.1 内存泄漏的幽灵ThreadLocal与上下文切换在压力测试运行几个小时后我们观察到JVM堆内存缓慢但持续地增长最终导致Full GC频繁甚至OOM。排查过程堆转储分析使用jmap和MAT工具分析堆转储发现大量ThreadLocal对象及其关联的上下文信息如MDC日志上下文、一些第三方库的缓存无法被回收。根因定位虚拟线程生命周期短暂创建和销毁频繁。一些库包括我们代码中使用了ThreadLocal来存储请求级别的上下文。在平台线程时代由于线程池复用线程ThreadLocal也能被复用。但在虚拟线程中一个虚拟线程结束后其载体线程可能立即被用来执行另一个完全不相关的虚拟线程而之前虚拟线程设置的ThreadLocal值如果没有被及时清理就会发生泄漏——因为载体线程被复用了但ThreadLocal里的旧数据还在。罪魁祸首我们发现了两个主要来源一是日志框架Logback的MDCMapped Diagnostic Context我们在拦截器中为每个请求设置了requestId二是一个内部使用的缓存工具类用了ThreadLocal来存储临时计算结果。解决方案使用ScopedValueJava 20预览Java 21正式替代ThreadLocalScopedValue是专为虚拟线程设计的提供了有界范围的、不可变的上下文传递。它在线程结束时会自动清理。private static final ScopedValueString REQUEST_ID ScopedValue.newInstance(); public void handleRequest(HttpServletRequest request) { String requestId generateId(); ScopedValue.where(REQUEST_ID, requestId).run(() - { // 在这个作用域内REQUEST_ID.get() 可以获取到 requestId process(); }); // 作用域结束值自动清理 }强制清理对于暂时无法替换的ThreadLocal使用场景如某些第三方库我们在虚拟线程执行结束前例如在Servlet Filter或Spring Interceptor的afterCompletion中主动调用ThreadLocal.remove()。升级依赖检查并升级所有第三方库到支持虚拟线程的最新版本许多主流框架如Log4j2、Micrometer的新版本都已修复了相关的ThreadLocal问题。4.2 锁与同步的“降级”风险虚拟线程鼓励使用阻塞IO但关于“锁”的行为需要特别注意。在平台线程上synchronized关键字或ReentrantLock锁住的是当前执行的线程。而在虚拟线程上锁住的是当前的虚拟线程。这听起来没问题但危险在于如果一个虚拟线程在持有锁的情况下执行了阻塞IO操作比如网络请求它会被挂起但锁仍然被它持有。如果这个IO操作耗时很长其他需要同一把锁的虚拟线程就会被长时间阻塞可能引发死锁或严重性能下降。案例我们有一个缓存加载器使用了“双重检查锁定”模式来懒加载一个热点配置。public class ConfigLoader { private volatile Config config; private final Object lock new Object(); public Config getConfig() { if (config null) { synchronized (lock) { // 虚拟线程A进入同步块 if (config null) { config loadConfigFromRemote(); // 这里进行网络IO虚拟线程A被挂起 } } } return config; } }当虚拟线程A在synchronized块内进行网络IO被挂起时锁未被释放。此时虚拟线程B调用getConfig()会在synchronized处被阻塞即使此时config仍是null。如果大量请求涌入所有虚拟线程都会阻塞在这个锁上服务瘫痪。解决方案原则尽量避免在持有锁的情况下执行任何可能阻塞的IO操作。重构对于上述缓存场景改用ConcurrentHashMap.computeIfAbsent或StampedLock等更灵活的并发工具或者将IO操作移到锁范围之外。例如先不加锁地检查如果为空再进行一个原子性的“获取-计算-存储”操作。使用ReentrantLock并显式控制如果必须用锁优先使用ReentrantLock因为它提供了更灵活的控制如可中断、可超时。但核心原则不变锁内不进行阻塞IO。4.3 监控与可观测性体系重建虚拟线程的引入使得传统的基于平台线程ID的监控链路如APM中的线程堆栈采样几乎失效。虚拟线程ID是连续递增的数字且生命周期短在日志和监控中直接打印线程ID意义不大。我们的监控改造链路追踪Tracing强化分布式链路追踪如使用OpenTelemetry。将traceId和spanId通过ScopedValue在虚拟线程之间传递确保整个调用链的上下文不丢失。这是监控虚拟线程应用最有效的手段。日志记录在日志模式中不再使用%thread而是记录ScopedValue中的requestId或链路追踪ID。同时可以记录虚拟线程的载体线程信息Thread.currentThread().toString()会包含载体线程信息用于深度调试。JVM指标关注新的JVM指标如jdk.VirtualThread.*如jdk.VirtualThread.count虚拟线程数jdk.VirtualThread.created创建总数。使用Micrometer等工具暴露这些指标到监控系统。线程转储传统的jstack或Thread.dumpAllStackTraces()对虚拟线程支持有限。需要使用jcmd pid Thread.dump_to_file -formatjson file来生成包含虚拟线程详细信息的转储文件进行分析。5. 效果评估与未来展望经过近两个月的重构、测试和灰度上线新系统最终全面取代了旧系统。性能对比吞吐量在相同的硬件资源下处理混合型IO轻度CPU工作负载的吞吐量提升了约40%。这主要得益于虚拟线程极低的创建和上下文切换开销使得我们能够用更少的资源支撑更高的并发连接数。资源利用率CPU利用率更加平稳避免了旧系统中因线程池排队或回调调度不均导致的CPU“毛刺”现象。内存方面在解决了ThreadLocal泄漏问题后内存增长曲线健康。延迟P99长尾延迟显著降低。旧系统中一旦线程池耗尽后续请求必须排队P99延迟会飙升。新系统中虚拟线程几乎可以“来一个请求就服务一个”排队现象大幅减少P99延迟降低了约60%。开发与维护效率这是最大的隐性收益。代码回归同步风格后可读性、可调试性大幅提升。新同事上手速度加快线上问题的平均排查时间MTTR缩短了超过50%。遇到的挑战与代价生态成熟度并非所有常用的Java库都完全适配了虚拟线程。我们需要对依赖库进行逐一评估、测试甚至打补丁或寻找替代方案。思维转变从“异步非阻塞”的响应式思维转变回“同步阻塞”但性能不差的思维需要团队一定的适应过程。尤其要时刻警惕“阻塞操作在锁内”等新的并发陷阱。调试工具如前所述调试和 profiling 工具链需要升级和适应。未来展望 这次重写验证了虚拟线程在IO密集型、高并发中间件系统如RAG中的巨大潜力。它在一定程度上简化了并发编程模型降低了心智负担。对于团队而言技术栈回归到更普适、更易理解的同步模型长期来看是利大于弊的。当然虚拟线程不是银弹对于计算密集型任务它并无优势甚至需要与平台线程池配合使用。下一步我们计划探索更多虚拟线程的高级特性例如更精细地使用StructuredTaskScope来管理复杂的子任务依赖树以及研究如何更好地与Project Loom的其他特性如Fiber结合。同时我们也会持续关注Java生态中主要组件对虚拟线程支持的进展逐步将最佳实践固化到团队的开发规范中。这次技术冒险虽然过程坎坷但结果令人满意它让我们对Java并发编程的未来有了更坚实的信心。
返回列表