Java构建流式RAG问答系统:架构分层与SSE实现详解 📅 2026/8/26 2:47:44 1. 项目概述从零到一构建一个流式响应的RAG问答系统最近在做一个内部知识库的智能问答功能核心需求很明确用户输入一个问题系统能快速从一堆文档里找到最相关的信息然后让大模型生成一个准确、流畅的回答并且最好能像ChatGPT那样答案一个字一个字地“流”出来提升用户体验。这不就是典型的RAG检索增强生成场景吗而且团队技术栈以Java为主所以很自然地我们就决定用纯Java技术栈来搭这个RAG问答的全链路。这个项目标题里的“架构分层”和“SSE流式”是两个关键点。架构分层意味着我们不能把所有代码都堆在一个Controller里那样后期维护和扩展会是噩梦。而SSE流式则是为了应对大模型生成回答可能比较慢的问题让前端能实时看到生成过程避免用户对着一个空白页面干等。整个链路拆开来看就是从用户提问开始经过问题理解、文档检索、答案生成、流式推送这几个核心环节。下面我就结合这次实战把每个环节的设计思路、技术选型、踩过的坑和最终方案详细拆解一遍。2. 核心架构分层设计清晰的责任边界是稳定的基石一开始最容易犯的错误就是“一把梭”把检索、调用模型、流式输出全写在一起。我们这次严格进行了分层核心思想是“高内聚、低耦合”让每一层只关心自己职责范围内的事情。2.1 分层模型与职责定义我们最终采用了经典的四层架构自顶向下分别是接口层、应用服务层、领域层和基础设施层。接口层Interface Layer这一层直接面向外部调用主要是Spring MVC的RestController。它的职责非常单一接收HTTP请求通常是用户的问题。进行基本的参数校验如问题不能为空。调用下一层的应用服务并处理其返回结果或异常。对于SSE流式接口它负责建立和维持SSE连接并将应用服务层返回的Flux响应式流数据通过SseEmitter推送给客户端。注意接口层绝对不应该包含任何业务逻辑比如“如何检索文档”、“怎么调用模型”这些都不归它管。它的代码应该非常“薄”。应用服务层Application Service Layer这是协调者。一个“问答”用例会涉及到检索文档、调用大模型等多个领域能力。应用服务层就负责把这些能力串起来编排成一个完整的业务流程。例如我们有一个QaService它的streamingAnswer方法内部大概会做这几件事调用检索服务获取相关文档片段 - 组装成大模型能理解的Prompt - 调用大模型服务并获取流式响应 - 对响应进行一些后处理如格式化。这一层的方法通常代表一个完整的用户用例。领域层Domain Layer这是业务核心包含了我们系统的核心概念和规则。在这一层我们定义了诸如Question问题、DocumentChunk文档片段、RelevantContext相关上下文、Answer答案等实体。同时这里也包含了核心的领域服务接口例如RetrievalService检索服务和LlmService大语言模型服务。这些接口定义了“要做什么”但不关心“具体怎么做”。领域层应该是技术无关的它不依赖任何具体的数据库、向量库或某个特定的AI模型SDK。基础设施层Infrastructure Layer这是具体实现的所在地。它实现了领域层定义的接口。例如RetrievalService的实现类会去操作Milvus或Elasticsearch执行向量相似度搜索或关键词搜索。LlmService的实现类会通过HTTP Client调用OpenAI、通义千问或本地部署的Ollama API。这一层还包含数据库访问JPA/MyBatis、文件解析、配置管理等“技术细节”。这样的分层带来了几个明显的好处一是代码结构清晰新人上手快二是便于测试可以轻松Mock基础设施层来测试应用服务三是未来要换向量库比如从Milvus换成PgVector或者换大模型从GPT换成DeepSeek只需要在基础设施层替换实现上层业务代码几乎不用动。2.2 技术选型背后的思考技术选型不是追新而是找最适合当前团队和场景的解决方案。Spring Boot Spring WebFlux基础框架没得说Spring Boot能快速搭起项目骨架。为什么引入WebFlux核心是为了SSE和更好的并发处理。传统的Servlet API处理SSE比较别扭而WebFlux的响应式编程模型和Flux类型与SSE是天作之合写起来非常自然。即使你现在不用SSE用Mono/Flux来封装一些异步或潜在的未来流式操作也能让代码更有弹性。向量数据库Milvus vs. Elasticsearch这是一个关键选择。我们评估了Milvus和ES的向量搜索插件。Milvus专为向量搜索设计性能指标尤其是高维、大规模向量和社区活跃度都很棒。如果你的场景是纯语义检索Milvus是首选。Elasticsearch dense_vector我们最终部分选择了这个方案。原因有二一是团队已有ES运维经验降低学习成本二是我们的需求是“混合检索”Hybrid Search即同时进行关键词匹配BM25和向量语义匹配然后将结果融合。ES原生支持BM25加上向量插件可以在一轮查询中完成混合检索架构更简单。如果选Milvus可能还需要单独维护一个ES或Lucene来做关键词检索再自己写融合逻辑复杂度更高。实操心得不要盲目追求“专业”向量库。评估你的核心检索模式。如果强相关性的关键词匹配很重要比如找产品型号、错误代码混合检索收益很大ES是更综合的选择。如果完全是开放域语义问答Milvus等专业向量库优势更明显。大模型接入OpenAI API vs. 本地模型我们采用了“双轨制”。对于生产环境调用云端API如GPT-4稳定、效果好。对于开发测试和成本敏感的内部场景我们集成了Ollama来本地运行诸如Qwen2.5、Llama3等开源模型。在Java中调用无非就是封装一个HTTP Client根据不同的模型提供商组装不同的请求体OpenAI格式、Ollama格式等。这里的关键是定义一个统一的LlmService接口将格式差异隐藏在具体实现中。SSE推送SseEmitter vs. 响应式流在Spring MVC环境下SseEmitter是标准选择使用简单。但在我们使用了WebFlux后更推荐直接返回ResponseBodyEmitter或Flux。特别是Flux它可以无缝对接下游的流式响应比如从大模型API返回的流实现从模型到客户端的“端到端”流式传输中间不需要缓冲整个回答内存效率更高。3. 检索问答全链路核心环节拆解链路的核心是“检索-生成”这个闭环。每一个环节的细节都直接影响最终答案的质量。3.1 文档处理与向量化入库这是RAG的“记忆”形成阶段离线进行但至关重要。流程是原始文档PDF/Word/TXT- 文本提取 - 文本分割切片- 向量化 - 存入向量库。文本分割Chunking这是第一个坑。不能简单按固定字符数比如512字切。那样很容易把一个完整的句子或概念从中间切断导致检索出来的片段语义不完整。我们采用了基于标记Token的滑动窗口法并尽量保证按段落、标题等自然边界切分。例如用langchain4j提供的DocumentSplitter工具可以设置chunkSize如1000 tokens和chunkOverlap如200 tokens。重叠部分保证了上下文连贯性避免信息丢失在切分边界。向量化模型选择我们测试了text-embedding-ada-002、bge-large-zh等模型。对于中文场景bge-large-zh表现确实不错。关键点是检索时的查询向量化模型必须与入库时使用的模型一致否则向量空间不一致相似度计算毫无意义。我们将模型调用也封装成服务保证线上线下模型统一。元数据存储除了向量我们还在ES里存储了片段的元数据如doc_id原文档ID、chunk_index片段序号、title所属标题等。这在后续的重排序和答案生成阶段很有用比如告诉模型“这个信息来源于XX文档的XX章节”。3.2 混合检索与重排序策略用户提问后系统如何找到最相关的文档片段我们实现了混合检索加精排的流程。1. 多路召回这是“广撒网”阶段目标是尽可能不错过任何可能相关的信息。向量检索路用问题的向量去向量库进行近似最近邻搜索ANN召回语义上最接近的Top K个片段比如K20。关键词检索路用问题中的关键词在ES中进行BM25检索召回相关性评分高的Top M个片段比如M20。注意事项两路召回的结果可能会有大量重复。需要根据片段ID进行去重。2. 分数归一化与融合两路检索的分数向量相似度分数和BM25分数量纲不同不能直接相加。我们采用了倒数排序融合Reciprocal Rank Fusion RRF。这种方法不关心具体分数值只关心排名。对于每个片段计算它在两路召回结果中的排名然后按照公式RRF score 1 / (rank k)分别计算k是一个常数通常取60最后将两路的RRF分数相加作为总得分。实践下来RRF比手动调权重的线性融合更鲁棒。3. 重排序Re-ranking经过融合我们得到了一个更全面的候选片段列表比如40个。但里面可能仍然存在一些语义相关但实际没用或者相关性不高的片段。这时需要“精挑细选”。我们引入了一个轻量级的交叉编码器Cross-Encoder模型如bge-reranker。它的原理是将“问题”和“每一个候选片段”同时输入模型进行深度的注意力交互从而得到一个更精细的相关性分数。虽然比向量检索慢但只对少量候选片段如40个进行开销可控。用重排序模型重新打分后选取Top N如N5作为最终提供给大模型的上下文。3.3 Prompt工程与上下文组装检索到相关片段后如何有效地“喂”给大模型Prompt设计是关键。我们不再使用简单的“请根据以下上下文回答问题”这种模板。一个更有效的Prompt结构如下你是一个专业的问答助手。请严格根据提供的“参考上下文”来回答问题。 如果上下文中的信息足以回答问题请基于上下文生成一个准确、简洁的答案。 如果上下文中的信息不足以完全回答问题你可以结合自己的知识进行补充但必须明确指出哪些信息来源于上下文哪些是你的补充。 如果上下文与问题完全无关请直接回答“根据提供的资料我无法回答这个问题”。 参考上下文 --- [文档片段1的标题] [文档片段1的内容] 来源文档1 --- [文档片段2的标题] [文档片段2的内容] 来源文档2 --- ...更多片段... 问题{用户输入的问题}关键点角色设定明确模型的身份引导其行为模式。指令清晰明确要求模型“基于上下文”并处理信息不足的情况。结构化上下文为每个片段添加标题和来源帮助模型更好地理解和引用。用---等分隔符清晰划分不同片段。位置放置将上下文放在问题之前符合模型的阅读习惯。在Java代码中我们使用StringBuilder或模板引擎如Thymeleaf、FreeMarker来动态组装这个Prompt字符串。确保片段内容被正确转义避免破坏Prompt结构。4. SSE流式输出完整实现流式输出是提升体验的利器。目标是让大模型生成一个词我们就立刻推给前端一个词。4.1 服务端实现WebFlux与响应式流在Spring WebFlux中实现一个SSE端点非常优雅。RestController RequestMapping(/api/rag) public class StreamQaController { Autowired private QaService qaService; GetMapping(value /stream-answer, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxServerSentEventString streamAnswer(RequestParam String question) { // 1. 调用应用服务层获取一个包含答案片段的Flux流 FluxString answerStream qaService.streamingAnswer(question); // 2. 将字符串流包装成SSE事件流 return answerStream .map(chunk - ServerSentEvent.Stringbuilder() .data(chunk) // 数据内容 .event(message) // 事件类型前端可监听 .id(UUID.randomUUID().toString()) // 可选事件ID .build()) .onErrorResume(e - { // 发生错误时发送一个错误事件 return Flux.just(ServerSentEvent.Stringbuilder() .event(error) .data(生成回答时发生错误: e.getMessage()) .build()); }) .doOnComplete(() - log.info(SSE流式问答完成问题: {}, question)); } }核心是QaService.streamingAnswer(String question)方法它返回一个FluxString。这个Flux的每个元素就是大模型实时生成的一个词或一个片段。这个方法的内部需要完成检索、Prompt组装并调用支持流式响应的大模型API。4.2 与大模型流式API的对接以调用OpenAI的流式Chat Completion API为例Service public class OpenAILlmService implements LlmService { private final WebClient webClient; private final ObjectMapper objectMapper; public FluxString streamCompletion(String prompt) { OpenAIChatRequest request buildRequest(prompt); // 构建请求设置stream: true return webClient.post() .uri(https://api.openai.com/v1/chat/completions) .header(HttpHeaders.AUTHORIZATION, Bearer apiKey) .contentType(MediaType.APPLICATION_JSON) .bodyValue(request) .retrieve() .bodyToFlux(String.class) // 关键响应体是SSE流每个事件是一个字符串 .takeUntil(s - s.contains([DONE])) // 遇到结束标志停止 .filter(s - s.startsWith(data: ) !s.contains([DONE])) .map(s - s.substring(6)) // 去掉 data: 前缀 .map(this::parseDeltaContent) // 解析JSON提取delta中的content .filter(content - content ! null !content.isEmpty()); } private String parseDeltaContent(String jsonLine) { try { JsonNode node objectMapper.readTree(jsonLine); JsonNode choices node.path(choices); if (choices.isArray() choices.size() 0) { JsonNode delta choices.get(0).path(delta); return delta.path(content).asText(); // 提取流式返回的内容增量 } } catch (Exception e) { log.warn(解析SSE数据行失败: {}, jsonLine, e); } return null; } }这里的关键是使用WebClient并以FluxString的形式接收响应体。OpenAI的流式API返回的是text/event-stream格式每一行是一个JSON对象。我们需要过滤掉非数据行和结束标志然后解析出每次返回的文本增量delta.content。对于本地Ollama其流式API接口类似返回格式可能是纯文本流或类似的JSON行流需要根据其文档调整解析逻辑。核心思想不变将HTTP响应流转换为一个FluxString。4.3 前端对接与用户体验优化前端使用EventSourceAPI来连接SSE端点。const eventSource new EventSource(/api/rag/stream-answer?question${encodeURIComponent(question)}); const answerDiv document.getElementById(answer); eventSource.addEventListener(message, function(event) { // 持续追加到页面 answerDiv.innerHTML event.data; }); eventSource.addEventListener(error, function(event) { console.error(SSE连接错误:, event); eventSource.close(); answerDiv.innerHTML br/span stylecolor:red连接中断或发生错误。/span; }); // 当用户离开页面或主动停止时关闭连接 window.onbeforeunload () eventSource.close();优化点连接管理前端需要在组件卸载或用户主动取消时调用eventSource.close()避免僵尸连接。错误处理监听error事件给用户友好的提示并重连或提供重试按钮。加载指示器在连接建立后、第一个数据块到达前可以显示“正在思考...”的动画。打字机效果为了更好的体验可以对接收到的每个字符或片段进行缓动动画输出模拟打字效果。5. 全链路集成与性能调优把各个模块拼装起来并让整个系统跑得又快又稳。5.1 服务编排与异步化在QaService中整个流程是异步编排的充分利用响应式编程的非阻塞特性。Service public class QaServiceImpl implements QaService { Override public FluxString streamingAnswer(String question) { // 1. 异步检索返回MonoListDocumentChunk MonoListDocumentChunk relevantChunksMono retrievalService.retrieveRelevantChunks(question); // 2. 组装Prompt等检索结果出来后同步组装 MonoString promptMono relevantChunksMono.map(chunks - buildPrompt(question, chunks)); // 3. 流式调用LLM将Prompt Mono转换为流式响应的Flux return promptMono.flatMapMany(prompt - llmService.streamCompletion(prompt)); } }这里flatMapMany是关键操作符它将一个包含Prompt的Mono“拍平”成LLM返回的Flux流。这样从检索结束到开始流式输出的延迟最小。5.2 关键性能指标与优化点首字延迟Time To First Token TTFT从用户发送问题到看到第一个字的时间。这是最重要的体验指标。优化TTFT检索优化确保向量索引性能考虑使用更快的嵌入模型如text-embedding-3-small。缓存对常见问题或高频查询的检索结果进行缓存。甚至可以对“问题向量”进行缓存避免重复计算。并行化如果混合检索的两路查询没有依赖可以并行执行使用Mono.zip或Flux.merge。生成速度Tokens Per Second TPS大模型输出token的速度。这主要取决于模型本身和网络延迟。对于本地模型确保服务器有足够的GPU资源。系统资源内存流式处理本身能极大降低内存峰值因为不需要缓冲完整响应。但仍需监控JVM堆内存特别是在处理大量并发SSE连接时。线程/连接数WebFlux默认使用Netty等非阻塞IO连接数上限很高。但需注意操作系统文件描述符限制和下游服务如ES、大模型API的并发连接限制合理配置连接池。5.3 可观测性与监控没有监控的系统就是在裸奔。我们集成了Micrometer和Prometheus来暴露指标。自定义指标rag.retrieval.duration检索阶段耗时直方图。rag.llm.ttft首字延迟直方图。rag.llm.total_duration整个生成耗时。rag.sse.connections.active当前活跃的SSE连接数。链路追踪在每个请求入口加入Trace ID并传递到检索服务、LLM调用等下游方便在日志中串联一次完整请求的路径快速定位瓶颈。日志记录结构化日志JSON格式记录每次问答的问题、检索到的文档ID、生成的答案可脱敏用于后续分析答案质量和优化检索效果。6. 常见问题排查与实战心得在实际开发和上线过程中我们遇到了不少问题这里总结一下。6.1 检索相关的问题问题1检索出来的片段似乎不相关。检查向量模型一致性确保入库和查询用的是同一个嵌入模型这是最常见的原因。调整切片策略如果片段总是支离破碎尝试调整chunkSize和chunkOverlap或者尝试按句子、段落分割。审视混合检索权重如果是混合检索可能是关键词检索权重太高淹没了语义检索的结果。尝试调整RRF中的参数k或者改用加权分数融合并调整权重。引入重排序如果候选片段列表里“似乎”有相关的但排名靠后引入重排序模型能极大改善最终结果。问题2检索速度慢。检查索引确保向量字段和文本字段都建立了合适的索引。对于ESdense_vector类型需要配置正确的相似度算法和索引选项如hnsw。限制召回数量在保证召回率的前提下尽量减少第一轮召回的top_k值比如从50降到20能显著减少重排序和后续处理的计算量。异步与缓存如5.2节所述使用异步编排和缓存。6.2 流式输出相关的问题问题1前端收不到流式数据或者连接很快断开。检查响应头服务端接口必须设置produces MediaType.TEXT_EVENT_STREAM_VALUE。检查网络代理和网关Nginx等代理服务器默认可能对响应进行缓冲buffer这会破坏流式传输。需要在Nginx配置中为对应路径关闭代理缓冲location /api/rag/stream-answer { proxy_pass http://backend; proxy_buffering off; # 关键 proxy_cache off; proxy_set_header Connection ; proxy_http_version 1.1; chunked_transfer_encoding off; # 对于SSE有时也需要关闭 }心跳保活SSE连接长时间没有数据发送可能会被代理或浏览器断开。可以在服务端定期发送注释行: heartbeat\n\n作为心跳。问题2流式输出中断或不完整。后端异常处理确保在Flux流中妥善处理异常使用onErrorResume返回一个错误事件而不是让整个流崩溃。参考4.1节的代码。检查LLM API稳定性如果是调用云端API网络波动或API限流可能导致流中断。需要实现重试机制但注意对于已经开始的流式响应重试逻辑比较复杂通常是在应用层记录失败提示用户重试。前端EventSource监听确保前端监听了error事件并能在UI上给予提示。6.3 答案质量相关的问题问题1答案胡编乱造幻觉。强化Prompt指令在Prompt中明确要求“严格基于上下文”并设定惩罚机制如“如果无法从上下文找到答案请说不知道”。提供更多上下文增加检索后提供给模型的片段数量Top N。但要注意上下文窗口长度限制。检查检索质量如果检索到的片段本身就不相关模型巧妇难为无米之炊。回溯到问题1去优化检索。问题2答案包含无关信息或格式混乱。后处理清洗在流式输出结束后或输出过程中可以加入简单的后处理逻辑比如过滤掉模型可能自生成的“根据以上上下文...”之类的套话或者用正则表达式规范一下答案的格式。Prompt中指定格式明确要求模型以某种格式输出例如“请用简洁的列表形式回答”、“请先给出结论再分点阐述”。踩过的一个大坑早期我们没有严格分离“检索”和“生成”的异常。一旦检索服务超时整个流式接口就卡住直到超时。后来我们将检索阶段设计为快速失败如果检索超时或失败立即向SSE流发送一个错误事件并结束流前端能立刻收到“检索失败”的提示而不是无限等待。这要求我们对retrievalService.retrieveRelevantChunks的Mono设置一个合理的超时时间timeout操作符并在超时后返回一个兜底的错误结果而不是让整个流挂起。整个项目做下来最大的体会是RAG系统是一个复杂的系统工程不仅仅是调个API。从文档预处理的质量到检索策略的精度再到Prompt的设计和流式传输的稳定性环环相扣。用Java构建这样的系统虽然不如Python生态有那么多现成的LangChain类库但通过清晰的分层设计和合理的组件选型完全可以打造出高性能、高可维护的生产级应用。特别是在应对复杂业务逻辑和需要高并发、稳定流式输出的场景下Java技术栈的优势就体现出来了。