Java RAG系统架构分层设计与SSE流式响应实战
1. 项目缘起:从单体混沌到分层清晰的RAG系统演进
去年年底,我们团队接手了一个内部知识库问答系统的重构任务。最初的版本是一个典型的“赶工”产物:所有的代码——从PDF解析、文本切片、向量化嵌入,到最后的检索与大模型生成——都挤在一个庞大的Spring Boot Controller里。每当业务方提出一个新的需求,比如想给不同的知识库配置不同的召回策略,或者想在回答时引用更精确的原文片段,我们都要在这个已经超过2000行的“上帝类”里小心翼翼地修改,生怕牵一发而动全身。更头疼的是前端交互,用户每次提问都要等待漫长的十几秒,页面一直转圈,体验极差。在一次关键的演示中,系统甚至因为处理一个稍复杂的查询而内存溢出,场面一度十分尴尬。
这次重构,我们的核心目标非常明确:第一, 架构分层 ,将混杂的职责清晰拆解,让系统变得可维护、可扩展;第二, 引入流式响应(SSE) ,彻底解决用户等待焦虑,实现答案的“边生成边返回”。这不仅仅是技术升级,更是对研发效率和用户体验的一次重要投资。经过几轮迭代,我们最终打磨出了一套稳定、高效且易于理解的Java RAG全链路实现方案。今天,我就把从架构设计到SSE流式推流的完整实战经验分享出来,希望能帮你绕过我们踩过的那些坑。
2. 核心架构分层设计:告别“意大利面条”代码
面对一个复杂的RAG系统,清晰的分层是保障其长期健康度的基石。我们摒弃了之前所有逻辑堆在一起的写法,采用了经典的分层架构思想,并针对RAG的特性进行了适配。
2.1 分层定义与职责边界
我们将系统自上而下划分为四个核心层,每一层都有其明确的单一职责。
表现层 (Presentation Layer) 这是系统对外的窗口,主要负责接收HTTP请求、解析参数、调用应用服务并返回响应。在这一层,我们重点关注的是 协议适配 。例如,对于普通的问答请求,我们提供同步的RESTful API;而对于需要流式输出的场景,则专门提供SSE(Server-Sent Events)协议的端点。这一层应该保持“薄”,只做协议转换和简单的参数校验,不包含任何业务逻辑。
应用层 (Application Layer) 这是协调整个RAG流程的“指挥中心”。它不关心数据具体从哪里来、向量怎么算,它只负责编排。一个典型的“问答”用例在此层的流程是:1. 接收用户问题;2. 调用检索服务获取相关文档片段;3. 将问题和文档片段组装成Prompt;4. 调用大模型服务生成答案;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处理PDF,Apache 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),关键词检索分数为 ks (BM25分数,归一化到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生成全部完成才能看到答案,对于长答案体验很差。SSE(Server-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. 执行检索(这部分通常较快,可以同步或异步)
List<Chunk> 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 stream=true)的流式响应本身也是一个数据流。我们需要建立一个管道,将LLM的流式输出“转发”为SSE事件。这里要注意 背压(Backpressure) 问题:如果LLM生成速度远快于网络发送速度,可能导致内存中积压大量数据。我们的 llmService.streamGenerate 方法内部使用回调机制,只有当SSE Emitter.send() 成功(意味着数据已被框架放入发送缓冲区),才会请求下一个LLM数据块,这是一种简单的客户端拉取式背压控制。
4.3 前端如何消费SSE流
前端使用 EventSource API可以轻松连接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 对可以并行的操作进行编排:
CompletableFuture<List<Chunk>> vectorFuture = CompletableFuture.supplyAsync(() -> vectorStore.similaritySearch(query), ioExecutor);
CompletableFuture<List<Chunk>> keywordFuture = CompletableFuture.supplyAsync(() -> esService.keywordSearch(query), ioExecutor);
CompletableFuture<List<Chunk>> mergedFuture = vectorFuture
.thenCombine(keywordFuture, this::mergeAndDeduplicate)
.thenApplyAsync(mergedList -> rerankService.rerank(query, mergedList), cpuExecutor);
这里需要注意 线程池隔离 :I/O操作(网络调用)使用一个较大的缓存线程池,CPU密集型操作(重排序计算)使用一个固定大小的线程池,避免相互影响。
5.3 内存管理与OOM预防
流式响应虽然改善了用户体验,但服务器端在生成完整答案前,可能需要将检索到的上下文和生成的文本都缓存在内存中。一次OOM(OutOfMemoryError)让我们印象深刻。
- 上下文长度限制 :严格限制送入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实战经验,能为你正在构建或优化的智能问答系统提供一些切实可行的思路和避坑参考。
更多推荐


所有评论(0)