ARTICLE DETAIL

资讯详情

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

Java RAG系统架构分层设计与SSE流式响应实战

Java RAG系统架构分层设计与SSE流式响应实战 1. 项目缘起从单体混沌到分层清晰的RAG系统演进去年年底我们团队接手了一个内部知识库问答系统的重构任务。最初的版本是一个典型的“赶工”产物所有的代码——从PDF解析、文本切片、向量化嵌入到最后的检索与大模型生成——都挤在一个庞大的Spring Boot Controller里。每当业务方提出一个新的需求比如想给不同的知识库配置不同的召回策略或者想在回答时引用更精确的原文片段我们都要在这个已经超过2000行的“上帝类”里小心翼翼地修改生怕牵一发而动全身。更头疼的是前端交互用户每次提问都要等待漫长的十几秒页面一直转圈体验极差。在一次关键的演示中系统甚至因为处理一个稍复杂的查询而内存溢出场面一度十分尴尬。这次重构我们的核心目标非常明确第一架构分层将混杂的职责清晰拆解让系统变得可维护、可扩展第二引入流式响应SSE彻底解决用户等待焦虑实现答案的“边生成边返回”。这不仅仅是技术升级更是对研发效率和用户体验的一次重要投资。经过几轮迭代我们最终打磨出了一套稳定、高效且易于理解的Java RAG全链路实现方案。今天我就把从架构设计到SSE流式推流的完整实战经验分享出来希望能帮你绕过我们踩过的那些坑。2. 核心架构分层设计告别“意大利面条”代码面对一个复杂的RAG系统清晰的分层是保障其长期健康度的基石。我们摒弃了之前所有逻辑堆在一起的写法采用了经典的分层架构思想并针对RAG的特性进行了适配。2.1 分层定义与职责边界我们将系统自上而下划分为四个核心层每一层都有其明确的单一职责。表现层 (Presentation Layer)这是系统对外的窗口主要负责接收HTTP请求、解析参数、调用应用服务并返回响应。在这一层我们重点关注的是协议适配。例如对于普通的问答请求我们提供同步的RESTful API而对于需要流式输出的场景则专门提供SSEServer-Sent Events协议的端点。这一层应该保持“薄”只做协议转换和简单的参数校验不包含任何业务逻辑。应用层 (Application Layer)这是协调整个RAG流程的“指挥中心”。它不关心数据具体从哪里来、向量怎么算它只负责编排。一个典型的“问答”用例在此层的流程是1. 接收用户问题2. 调用检索服务获取相关文档片段3. 将问题和文档片段组装成Prompt4. 调用大模型服务生成答案5. 处理并返回结果。这一层是**用例Use Case**的体现每个公开的API背后通常对应一个应用服务方法。领域层 (Domain Layer)这是系统的核心和灵魂包含了RAG业务的核心概念与规则。我们在这里定义了诸如Document文档、Chunk文本切片、Embedding向量、RetrievalResult检索结果、LLMResponse大模型响应等实体Entity和值对象Value Object。同时一些核心的业务逻辑比如文本切片策略是按段落切还是按固定长度切、检索结果的融合与去重规则等也会以领域服务Domain Service的形式放在这一层。这一层应该是技术无关的即不依赖任何特定的框架如Spring或外部库如某个向量数据库的客户端。**基础设施层 (Infrastructure Layer) **这是所有技术细节的“实现区”负责为上层提供具体的技术能力。它主要包括向量数据库客户端实现向量的存储、检索相似度搜索等操作可能是对Pinecone、Milvus、Elasticsearch带向量插件或PGVector等库的封装。大模型客户端封装调用OpenAI API、通义千问、DeepSeek等大模型服务的细节处理认证、请求构造和响应解析。文件解析器集成Apache PDFBox、Apache Tika等库实现PDF、Word、TXT等不同格式文件的文本提取。嵌入模型客户端调用如OpenAI的text-embedding-ada-002、BGE、M3E等模型将文本转换为向量。持久化存储使用JPA、MyBatis等操作关系型数据库存储知识库元数据、操作日志等。分层之后依赖关系变得清晰且单向表现层依赖应用层应用层依赖领域层而领域层则依赖基础设施层提供的接口Interface。这种依赖倒置使得我们可以轻松替换底层实现例如将向量数据库从Milvus切换到Weaviate只需在基础设施层提供新的实现上层业务代码几乎无需改动。2.2 分层后的代码组织与包结构清晰的架构也需要清晰的代码组织来体现。我们的项目包结构大致如下src/main/java/com/yourcompany/rag/ ├── application/ # 应用层 │ ├── service/ # 应用服务如 QaApplicationService │ └── dto/ # 入参出参对象 ├── domain/ # 领域层 │ ├── model/ # 实体与值对象 │ ├── service/ # 领域服务如 ChunkingService, RerankService │ └── repository/ # 仓储接口由基础设施层实现 ├── infrastructure/ # 基础设施层 │ ├── persistence/ # 持久化实现JPA等 │ ├── vectorstore/ # 向量数据库客户端封装 │ ├── llm/ # 大模型客户端封装 │ └── parser/ # 文件解析器 └── presentation/ # 表现层 ├── controller/ # 控制器包含REST和SSE端点 └── web/ # Web相关配置如SSE的Emitter管理这样的结构让新成员能快速理解系统脉络定位功能代码也变得异常轻松。3. 检索问答全链路核心环节拆解有了清晰的分层架构作为骨架我们来填充RAG最核心的“检索-生成”链路上的血肉。这个过程可以细化为多个步骤每一步的选择都直接影响最终效果。3.1 知识入库从文档到向量在问答之前我们需要先构建知识库。这个过程是离线的但设计好坏决定了检索质量的上限。文档解析与清洗我们使用Apache PDFBox处理PDFApache Tika作为格式探测和后备解析器。解析出的原始文本往往包含大量噪音无意义的页眉页脚、重复的换行符、乱码字符等。我们实现了一个TextCleaner组件通过正则表达式和启发式规则进行清洗。例如连续超过3个换行符替换为1个过滤掉只包含页码或“保密”字样的行。这一步看似琐碎却能显著提升后续切片和嵌入的质量。文本切片Chunking策略这是RAG的“阿喀琉斯之踵”。切得太碎上下文信息丢失切得太大会引入无关噪声且影响嵌入和检索效率。我们采用了分层切片的策略语义切片优先尝试按自然段落\n\n或Markdown/PDF的标题进行切分。这能最好地保留语义完整性。固定长度重叠切片对于长段落或无结构文本采用滑动窗口。我们设置窗口大小为500字符重叠为50字符。这确保了上下文连续性避免在窗口边界切断关键信息。注意重叠不是越大越好。过大的重叠如50%会急剧增加向量存储和检索的计算开销可能带来边际收益递减。通常10%-20%的重叠是一个不错的起点。向量化嵌入与存储清洗和切片后的文本通过嵌入模型转换为向量。我们封装了一个EmbeddingService内部可以适配不同的嵌入模型API。这里的关键是异步批量处理。同步地一条条请求嵌入API是性能瓶颈。我们使用CompletableFuture或Project Reactor将文本批量例如每100条一批发送并设置合理的超时和重试机制。生成向量后连同原文片段chunk、所属文档ID、元数据如切片索引一并存入向量数据库。我们选择PGVector因为团队PostgreSQL熟其vector类型和-余弦距离操作符用起来非常直观。3.2 在线检索多路召回与重排序当用户提问时系统进入在线检索阶段。我们实践了“多路召回重排序”的范式来提升召回质量。查询理解与向量召回首先对用户原始查询query进行轻量级处理如纠错、同义词扩展使用WordNet或业务词表生成优化后的查询query_opt。将其向量化后在向量数据库中进行相似度搜索KNN。这是最核心的向量召回路径。关键词召回作为补充单纯依赖向量检索有时会错过那些关键词匹配度高但语义表达不同的文档。因此我们并行地走关键词召回BM25路径。我们将文档切片也同步索引到Elasticsearch中使用BM25算法进行全文检索。这一步召回的是与查询词直接匹配的片段。多路结果融合与去重两路召回会各自返回一个Top-K的列表比如各20条。直接合并会面临两个问题1. 结果有重叠2. 如何排序我们采用简单的加权分数融合策略。假设向量检索分数为vs余弦相似度0-1关键词检索分数为ksBM25分数归一化到0-1则融合分数fs α * vs (1-α) * ks。α是一个可调参数我们根据业务测试设为0.7更偏向语义相似度。然后根据fs对合并后的列表重新排序并基于切片ID或内容哈希进行去重。重排序Reranking精炼经过融合排序后的列表例如30条已经不错但还可以用更精细但更耗资源的重排序模型再精炼一次。我们引入了一个轻量级的交叉编码器Cross-Encoder例如BGE-Reranker。它将查询和每个候选文档片段一起输入模型直接计算一个相关性分数。这个分数比嵌入向量的余弦相似度更精准。我们对Top-N如10条的融合结果进行重排序并替换其分数。这一步虽然增加了几十到几百毫秒的延迟但对于最终答案的质量提升尤其是在需要精确匹配的场景下效果显著。3.3 提示工程与大模型调用检索到最相关的几个文档片段后需要将它们和问题一起交给大模型生成最终答案。这里的核心是构造一个清晰的Prompt。Prompt模板设计我们避免将原始片段简单拼接。一个结构化的Prompt模板如下你是一个专业的助手请严格根据以下提供的上下文信息来回答问题。如果上下文信息不足以回答问题请直接说“根据已知信息无法回答该问题”不要编造信息。 上下文信息 {context} 问题{question} 请根据上下文回答这里的{context}就是我们检索并重排序后的Top片段用\n\n---\n\n等分隔符连接。我们严格控制上下文的长度避免超过模型令牌限制。大模型客户端封装我们封装了一个LlmService内部可以灵活切换不同的模型提供商OpenAI, Azure OpenAI, 国内各大模型API。关键点包括连接池与超时使用HTTP客户端连接池如Apache HttpClient或OkHttp管理连接设置连接、读写超时如30秒。故障转移与降级当主用模型API不可用时具备快速切换到备用模型的能力。流式支持对于需要流式返回的场景该服务需要能够处理分块chunked的响应这是实现SSE的基础。4. SSE流式输出实战让答案“流”起来同步请求下用户需要等待检索、LLM生成全部完成才能看到答案对于长答案体验很差。SSEServer-Sent Events协议允许服务器主动向浏览器推送数据是实现流式输出的理想选择。4.1 SSE协议简介与Spring Boot集成SSE是一种基于HTTP的轻量级协议服务器通过Content-Type: text/event-stream的响应可以持续发送多个data:事件。Spring Framework从4.2版本开始就提供了对SSE的原生支持核心是SseEmitter类。在Spring Boot中创建一个流式端点非常简单RestController RequestMapping(/api/rag) public class StreamQaController { GetMapping(value /stream-ask, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter streamAsk(RequestParam String question) { // 设置超时时间例如5分钟 SseEmitter emitter new SseEmitter(5 * 60 * 1000L); // 异步处理避免阻塞当前线程 CompletableFuture.runAsync(() - { try { // 1. 执行检索这部分通常较快可以同步或异步 ListChunk relevantChunks retrievalService.retrieve(question); // 2. 构造Prompt String prompt promptBuilder.build(question, relevantChunks); // 3. 调用支持流式响应的LLM服务 llmService.streamGenerate(prompt, new StreamCallback() { Override public void onData(String chunk) { try { // 将LLM返回的每一个文本块作为SSE事件发送 emitter.send(SseEmitter.event().data(chunk)); } catch (IOException e) { // 处理发送失败可能是客户端已断开 emitter.completeWithError(e); } } Override public void onComplete() { // 流式生成结束发送完成事件 emitter.complete(); } Override public void onError(Throwable t) { emitter.completeWithError(t); } }); } catch (Exception e) { emitter.completeWithError(e); } }); // 设置Emitter结束时的回调用于资源清理 emitter.onCompletion(() - log.info(SSE stream completed.)); emitter.onTimeout(() - log.warn(SSE stream timed out.)); emitter.onError((ex) - log.error(SSE stream error., ex)); return emitter; } }4.2 处理LLM流式响应与背压大模型API如OpenAI的Chat Completion with streamtrue的流式响应本身也是一个数据流。我们需要建立一个管道将LLM的流式输出“转发”为SSE事件。这里要注意背压Backpressure问题如果LLM生成速度远快于网络发送速度可能导致内存中积压大量数据。我们的llmService.streamGenerate方法内部使用回调机制只有当SSEEmitter.send()成功意味着数据已被框架放入发送缓冲区才会请求下一个LLM数据块这是一种简单的客户端拉取式背压控制。4.3 前端如何消费SSE流前端使用EventSourceAPI可以轻松连接SSE端点。const eventSource new EventSource(/api/rag/stream-ask?question encodeURIComponent(question)); const answerDiv document.getElementById(answer); eventSource.onmessage (event) { // 不断追加收到的数据块 answerDiv.innerHTML event.data; }; eventSource.onerror (error) { console.error(SSE error:, error); eventSource.close(); // 显示错误信息 };这样用户就能看到答案一个字一个字“打出来”的效果体验流畅度大幅提升。4.4 流式场景下的错误处理与连接管理流式连接生命周期长稳定性挑战更大。客户端断开用户关闭页面或刷新SseEmitter.send()会抛出IOException。我们需要在回调中捕获并优雅地终止后续的LLM调用如果可能避免服务器资源浪费。服务器错误在流式生成过程中如果发生错误如LLM API调用失败我们通过emitter.completeWithError()发送一个错误事件。前端可以监听EventSource的onerror事件给用户友好提示。连接超时设置合理的SseEmitter超时时间如5分钟并监听onTimeout进行清理。对于超长文本生成可以考虑“心跳”机制定期发送注释事件保持连接。Emitter管理在高并发下需要管理大量并发的SseEmitter实例防止内存泄漏。Spring Boot默认能处理但在极端情况下可以考虑用一个ConcurrentHashMap来跟踪活跃的Emitter并在onCompletion/onTimeout回调中将其移除。5. 性能优化与踩坑实录在实战中我们遇到了不少性能瓶颈和意料之外的问题以下是部分总结。5.1 向量检索的性能陷阱最初我们直接对千万级别的向量进行全量余弦相似度计算响应时间在秒级。优化方案索引是关键PGVector支持ivfflat或hnsw索引。我们在向量字段上创建了hnsw索引召回速度提升了一个数量级。创建索引的命令类似CREATE INDEX ON chunks USING hnsw (embedding vector_cosine_ops);。限制召回数量与精度检索时使用ORDER BY embedding - query_vector LIMIT K并合理设置K值如100。对于ivfflat索引还可以通过SET ivfflat.probes 10;来平衡速度与精度。连接池使用HikariCP等连接池管理数据库连接避免每次检索都新建连接。5.2 异步编排与资源竞争整个RAG链路涉及多个I/O密集型操作向量DB检索、关键词DB检索、重排序模型调用、LLM调用。如果全部同步串行延迟会叠加。我们使用CompletableFuture对可以并行的操作进行编排CompletableFutureListChunk vectorFuture CompletableFuture.supplyAsync(() - vectorStore.similaritySearch(query), ioExecutor); CompletableFutureListChunk keywordFuture CompletableFuture.supplyAsync(() - esService.keywordSearch(query), ioExecutor); CompletableFutureListChunk mergedFuture vectorFuture .thenCombine(keywordFuture, this::mergeAndDeduplicate) .thenApplyAsync(mergedList - rerankService.rerank(query, mergedList), cpuExecutor);这里需要注意线程池隔离I/O操作网络调用使用一个较大的缓存线程池CPU密集型操作重排序计算使用一个固定大小的线程池避免相互影响。5.3 内存管理与OOM预防流式响应虽然改善了用户体验但服务器端在生成完整答案前可能需要将检索到的上下文和生成的文本都缓存在内存中。一次OOMOutOfMemoryError让我们印象深刻。上下文长度限制严格限制送入LLM的上下文总长度。我们会优先选择相关性分数最高的片段直到总令牌数接近模型上限如16K的80%即停止。流式输出的缓冲区Spring的SseEmitter和底层Servlet容器会有输出缓冲区。如果LLM生成极快而网络极慢缓冲区可能积压。除了背压控制还可以考虑在应用层设置一个小的阻塞队列。JVM参数调优针对频繁创建大量临时对象如字符串、DTO的场景适当调年轻代-Xmn大小并选择适合的GC算法如G1。5.4 监控、日志与可观测性一个线上系统没有监控就是“盲人骑瞎马”。我们做了以下工作关键指标埋点使用Micrometer对每个环节耗时解析、切片、向量化、检索、LLM生成进行计时并统计QPS、错误率。分布式链路追踪集成SkyWalking或Zipkin为每个用户请求分配一个Trace ID贯穿所有微服务或异步线程方便定位慢查询。结构化日志使用JSON格式输出日志包含请求ID、用户ID、检索到的文档ID列表、LLM输入输出的令牌数等。这对分析效果、复现问题至关重要。LLM输入输出采样在日志中按一定比例采样记录完整的Prompt和生成的Answer用于人工评估效果和优化Prompt注意要对敏感信息进行脱敏。6. 总结与展望通过这次从混沌单体到清晰分层的重构并结合SSE流式输出我们的RAG系统在可维护性、响应速度和用户体验上都获得了质的飞跃。架构分层让团队协作效率提升新人也能快速上手流式输出则让终端用户感受到了实时交互的畅快。回顾整个过程我认为有几个决策点尤为关键第一领域层的抽象要足够纯粹它是对业务本质的描述不应被技术细节污染第二异步化与流式化需要从头设计而不是事后打补丁这关系到整个链路的资源模型第三监控和可观测性必须与功能开发同步进行否则线上问题排查将异常痛苦。未来我们计划在几个方向继续探索一是引入更智能的查询路由Query Routing根据问题类型决定是走RAG、直接调用LLM知识还是查结构化数据库二是尝试自愈Self-Correction机制让LLM对自身基于上下文的回答进行可信度评估并在置信度低时自动调整检索策略或提示用户三是探索Graph RAG利用知识图谱来建模文档间的关系实现更深层次的推理和问答。技术架构没有银弹最适合当前团队和业务场景的才是最好的。希望我们这套基于Java生态的RAG实战经验能为你正在构建或优化的智能问答系统提供一些切实可行的思路和避坑参考。
返回列表