ARTICLE DETAIL

资讯详情

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

RAG检索主链路实战:LangGraph图结构、混合召回与SSE流式输出

RAG检索主链路实战:LangGraph图结构、混合召回与SSE流式输出 1. 检索主链路到底在解决什么问题很多人做RAG项目前面文档解析、切块、向量化都跑通了一到用户提问→系统回答这条链路就卡住。要么是检索回来的内容跟问题八竿子打不着要么是拼出来的Prompt超了模型上下文窗口要么是前端等半天一个字都不吐。这一章要干的事就是把这条链路从头到尾串起来让它真正跑通一次完整的问答闭环。所谓检索主链路说白了就是四个动作的串联用户提问进来 → 检索相关文档片段 → 组装上下文喂给模型 → 流式返回答案。听起来简单但每一步都有坑。检索环节涉及向量召回和关键词召回的配合上下文组装涉及预算分配和父窗口回填流式返回涉及SSE协议和断连处理。这四个环节任何一个出问题整个闭环就断了。我见过太多项目检索用的是最朴素的向量相似度Top-K结果用户问上季度的销售数据召回来的全是销售部门组织架构这种语义相近但答非所问的片段。也见过上下文直接无脑拼接把检索到的十个片段全塞进去结果超出模型窗口要么报错要么被截断答案缺胳膊少腿。还有SSE流式接口写完本地测试好好的一上生产环境就频繁断连前端收到一半就没了。这一章的目标很明确让第一次问答闭环真正跑起来并且跑得稳。我会把LangGraph的图结构怎么设计、检索节点怎么写、上下文预算怎么算、父窗口回填怎么做、SSE流式接口怎么封装全部拆开讲清楚。每个决策背后的理由我都会说明每个参数我都会给出计算依据。你跟着走一遍就能得到一个可运行、可调试、可扩展的检索主链路。适合谁看如果你已经完成了文档入库和向量化正准备把问答接口串起来这一章就是为你写的。如果你还在纠结切块策略建议先回去把上一章的内容消化掉因为检索主链路的效果很大程度上取决于切块质量。如果你已经有一套能跑的RAG系统但效果不稳定、流式输出老出问题这一章里的排查思路和参数调优经验同样适用。2. LangGraph图结构的设计取舍2.1 为什么不用简单的函数调用链最朴素的RAG实现就是一个函数串到底retrieve(query) → build_prompt(docs, query) → call_llm(prompt) → return。这种写法在Demo阶段没问题但一旦要加多路召回、重排序、上下文压缩、多轮对话状态管理函数链就会变成一团乱麻。你没法在中途插入条件分支没法做并行检索没法在某个节点失败时优雅降级。LangGraph的价值在于把这条链路显式地建模成一张有向图。每个节点是一个处理单元每条边定义了流转方向状态在节点之间传递。这样做的好处是检索和生成解耦了你可以在检索节点后面挂一个重排序节点也可以并行跑向量召回和关键词召回再合并还可以在生成节点前加一个上下文预算检查节点超预算就触发压缩或截断。我实际项目里的图结构是这样的入口节点接收用户query然后分两条边并行走到向量检索节点和关键词检索节点两个节点的结果汇合到融合排序节点排序后的文档进入上下文组装节点组装时做预算检查和父窗口回填最后进入生成节点调用LLM并流式输出。整张图的状态用一个TypedDict管理包含query、vector_docs、keyword_docs、merged_docs、context、answer等字段。注意LangGraph的节点函数必须是纯函数或者只通过状态通信不要在节点里直接操作全局变量。我踩过这个坑多个请求并发时状态串了排查了半天才发现是节点里用了模块级变量。2.2 状态字段的设计与并发安全状态设计是LangGraph里最容易埋雷的地方。我的建议是状态字段尽量扁平避免嵌套过深。比如不要设计成{retrieval: {vector: [...], keyword: [...]}}而是直接{vector_docs: [...], keyword_docs: [...]}。扁平结构在节点间传递时更清晰调试时打印状态也一目了然。并发安全方面LangGraph本身对状态更新是做了隔离的每个请求有独立的状态副本。但如果你在节点里调用了外部服务比如向量数据库客户端要确保这个客户端是线程安全的或者每个请求独立创建。我用的Milvus客户端就是全局单例但它的查询方法是线程安全的所以没问题。如果你用的是某些非线程安全的SDK建议在节点内部按需创建连接。另一个容易忽略的点是状态字段的默认值。LangGraph在初始化状态时如果你没给某个字段赋值它可能是None。节点函数里如果直接对这个字段做操作就会报NoneType错误。我的做法是在图的入口节点里把所有字段都初始化一遍该给空列表的给空列表该给空字符串的给空字符串。from typing import TypedDict, List class RAGState(TypedDict): query: str vector_docs: List[dict] keyword_docs: List[dict] merged_docs: List[dict] context: str answer: str token_budget: int used_tokens: int这个状态定义里token_budget和used_tokens是给上下文预算控制用的后面会详细讲。2.3 条件边与降级路径图结构里必须考虑失败降级。向量检索超时了怎么办关键词检索返回空了怎么办融合排序后文档数量不够怎么办这些都需要条件边来处理。我的做法是在向量检索节点后面加一个条件判断如果返回结果为空或超时就走一条边到仅关键词检索的降级路径如果两个检索都失败就走无检索直接生成的兜底路径同时给用户返回一个提示说明当前检索服务不可用答案可能不准确。条件边的实现用LangGraph的add_conditional_edges方法传入一个判断函数根据状态里的字段决定走哪条边。判断函数要尽量简单只做路由决策不要在里面做复杂计算。def route_after_vector_search(state: RAGState) - str: if not state[vector_docs]: return keyword_only return merge graph.add_conditional_edges( vector_search, route_after_vector_search, {keyword_only: keyword_search, merge: merge_node} )这种设计的好处是即使某一路检索挂了整个问答链路还能继续跑用户体验不会完全中断。生产环境里可用性比完美答案更重要。3. 向量召回与关键词召回的配合逻辑3.1 纯向量检索的盲区在哪里向量检索擅长语义匹配用户问怎么提升团队效率它能召回团队协作方法效率工具推荐这类语义相近的文档。但它有个致命盲区对精确匹配不敏感。用户问OKR和KPI的区别向量检索可能召回一堆绩效管理的文档但真正讲OKR和KPI对比的那篇可能排在后面。另一个盲区是专有名词和缩写。比如用户问SSE接口怎么封装向量模型可能把SSE理解成服务器发送事件的语义但如果你文档里SSE是作为技术缩写出现的向量检索的召回率会明显下降。还有数字、日期、版本号这类信息向量检索基本无能为力。我实测过一个案例知识库里有篇文档标题是Ch07 检索主链路用户问第七章讲什么纯向量检索召回的却是主链路设计原则这种语义相近但章节不对的文档。后来加了关键词检索用Ch07做精确匹配才把正确的文档捞回来。3.2 关键词检索的补位策略关键词检索我用的是BM25算法配合jieba分词做中文处理。BM25的核心思想是一个词在文档中出现频率越高、在整个语料中出现频率越低它的权重就越大。这正好补上了向量检索对精确匹配不敏感的短板。具体实现上我用Elasticsearch做关键词检索的存储和查询。文档入库时除了存向量也把原文和分词后的term存进去。查询时对用户query做同样的分词然后用BM25算分。这里有个细节中文分词要考虑领域词典。比如检索主链路这个词通用分词器可能切成检索/主/链路但如果你把检索主链路加到自定义词典里它就会作为一个整体term匹配精度会高很多。import jieba jieba.add_word(检索主链路) jieba.add_word(父窗口回填) jieba.add_word(上下文预算) def tokenize(text: str) - List[str]: return list(jieba.cut(text))自定义词典这个事我建议在项目初期就建起来把领域内的专有名词、产品名、技术缩写都加进去。后期文档多了再补成本会高很多。3.3 两路结果的融合排序向量检索返回的是余弦相似度分数关键词检索返回的是BM25分数这两个分数不在一个量纲上不能直接相加。我的做法是先各自归一化再加权融合。归一化用Min-Max归一化把每路检索的分数映射到0到1之间。然后给两路分配权重我默认是向量0.7、关键词0.3。这个权重不是拍脑袋定的是根据实际效果调的。如果你的场景里用户提问偏口语化、语义化向量权重要高一些如果偏精确查询、术语查询关键词权重要高一些。融合的时候还要考虑去重。同一篇文档可能同时被两路检索召回这时候要合并成一条分数取加权后的最大值或者平均值。我用的是加权求和后再除以出现次数避免同一文档因为被两路都召回而分数虚高。融合策略适用场景优点缺点加权求和通用场景实现简单可调权重需要归一化权重难调RRF倒数排序多路召回无需归一化鲁棒性好丢失分数信息最大值融合精确匹配优先保留最强信号可能忽略互补信息交叉编码重排精度要求高精度最高计算开销大我实际用的是加权求和然后在融合排序节点后面加了一个可选的交叉编码重排节点。重排用的小模型是bge-reranker-base对Top-20的文档做精排取Top-5进入上下文组装。这一步能把最终答案的准确率提升10%到15%代价是增加约200毫秒的延迟。如果你的场景对延迟敏感可以跳过重排直接取融合后的Top-5。4. 上下文预算与父窗口回填的实操细节4.1 上下文预算到底怎么算上下文预算这个词听起来很玄其实就是一个简单的算术题模型的最大上下文窗口减去系统提示词占用的token数再减去预留的输出token数剩下的就是检索文档能用的token数。以GPT-4o为例最大上下文窗口是128K token。系统提示词我写的是约500 token预留输出2000 token那么检索文档可用的预算是128000 - 500 - 2000 125500 token。但实际使用中我不会把预算用满因为token计数有误差而且留一些余量能让模型有更多空间做推理。我通常把实际使用量控制在预算的80%左右。token计数用tiktoken库这是OpenAI官方的计数工具准确度很高。但要注意不同模型的tokenizer不一样如果你用的是Claude或者国产模型要用对应的计数工具。我项目里做了个抽象层根据模型名称选择对应的计数器。import tiktoken def count_tokens(text: str, model: str gpt-4o) - int: encoding tiktoken.encoding_for_model(model) return len(encoding.encode(text))预算分配上我的策略是按文档相关性分数从高到低依次加入上下文直到预算用尽。每加入一篇文档就计算一次累计token数超过预算就停止。这样能保证最相关的文档一定在上下文里。4.2 父窗口回填解决碎片化问题切块策略有个天然矛盾块切得小检索精度高但上下文碎片化块切得大上下文完整但检索精度下降。父窗口回填就是解决这个矛盾的。具体做法是检索时用小块做向量匹配命中后不直接返回小块而是返回这个小块所属的父块通常是更大的段落或整节内容。这样既保证了检索精度又保证了上下文的完整性。实现上我在文档入库时给每个小块记录一个parent_id父块单独存储。检索命中小块后根据parent_id去取父块内容。父块的大小我控制在1000到1500 token之间太小了上下文不够太大了浪费预算。提示父窗口回填会增加一次数据库查询如果父块和子块存在同一个集合里查询会很快如果分集合存储要注意加索引。我一开始没加索引回填时查一次要200毫秒后来在parent_id上建了索引降到5毫秒以内。回填之后还要做去重。多个小块可能属于同一个父块回填后会出现重复的父块内容。我的做法是用父块的ID做去重同一个父块只保留一次分数取该父块下所有命中子块中的最高分。4.3 预算超限时的截断与压缩即使做了预算控制有时候单篇父块就超了预算或者多篇父块加起来超了。这时候需要截断或压缩。截断的策略是从文档末尾开始截因为大部分文档的核心信息在前面。但这不是绝对的有些文档结论在后面截断会丢关键信息。我的做法是给每篇文档算一个信息密度分数密度高的文档优先保留完整密度低的文档从末尾截断。压缩的策略是用一个小模型对超长文档做摘要把摘要放进上下文。这个方案效果好但增加延迟和成本。我一般只在预算超限严重超过20%时才启用压缩轻微超限直接截断。def truncate_to_budget(docs: List[dict], budget: int) - List[dict]: result [] used 0 for doc in sorted(docs, keylambda x: x[score], reverseTrue): doc_tokens count_tokens(doc[content]) if used doc_tokens budget: result.append(doc) used doc_tokens else: remaining budget - used if remaining 100: doc[content] truncate_text(doc[content], remaining) result.append(doc) break return result这段代码里有个细节remaining 100这个判断意思是如果剩余预算不到100 token就不加了因为太短的片段对答案没帮助反而可能干扰模型。5. SSE流式接口的封装与断连处理5.1 为什么选SSE而不是WebSocket流式输出方案有SSE和WebSocket两种。SSE是单向的服务器推送WebSocket是双向通信。RAG问答场景里客户端只需要接收服务器推送的token流不需要向服务器发消息所以SSE足够了。SSE的优势是基于HTTP协议不需要额外的协议升级兼容性好实现简单。浏览器原生支持EventSource后端用FastAPI的StreamingResponse就能实现。WebSocket虽然功能更强但要处理连接管理、心跳、重连复杂度高很多。我项目里用的是FastAPI SSE。后端定义一个/chat/stream接口返回StreamingResponse媒体类型是text/event-stream。每次LLM生成一个token就通过SSE推送给前端。from fastapi import FastAPI from fastapi.responses import StreamingResponse app FastAPI() async def event_generator(query: str): async for token in rag_pipeline.astream(query): yield fdata: {json.dumps({token: token})}\n\n yield data: [DONE]\n\n app.get(/chat/stream) async def chat_stream(query: str): return StreamingResponse( event_generator(query), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no } )这里有几个关键点Cache-Control: no-cache防止中间层缓存X-Accel-Buffering: no防止Nginx缓冲这个坑我踩过不加这个头Nginx会把整个响应缓冲完再一次性发给客户端流式效果就没了[DONE]标记用来告诉前端流结束了。5.2 前端EventSource的封装前端用EventSource接收SSE流。但EventSource有个限制只支持GET请求不支持POST。如果你的query很长或者需要传复杂的参数GET的URL长度可能不够。解决方案是用fetch ReadableStream手动解析SSE流。async function streamChat(query, onToken, onDone) { const response await fetch(/chat/stream?query encodeURIComponent(query)); const reader response.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const lines buffer.split(\n\n); buffer lines.pop(); for (const line of lines) { if (line.startsWith(data: )) { const data line.slice(6); if (data [DONE]) { onDone(); return; } try { const parsed JSON.parse(data); onToken(parsed.token); } catch (e) { console.warn(解析SSE数据失败, e); } } } } }这段代码里buffer的作用是处理跨chunk的数据。SSE的消息以\n\n分隔但网络传输时一个消息可能被拆到两个chunk里。用buffer暂存不完整的部分等下一个chunk到了再拼接这是处理流式数据的标准做法。5.3 断连排查与超时配置SSE最常见的报错就是stream disconnected before completion: idle timeout waiting for sse。这个报错的意思是连接空闲超时服务器或中间层主动断开了连接。排查思路分三层客户端、中间层、服务端。客户端层面检查EventSource或fetch有没有设置超时。浏览器默认的fetch没有超时但如果你用了axios之类的库默认超时可能是60秒。LLM生成慢的时候60秒可能不够。中间层层面Nginx默认的proxy_read_timeout是60秒超过这个时间没有数据传输就会断开。解决方案是在Nginx配置里把这个值调大比如300秒。同时要确保proxy_buffering off否则Nginx会缓冲响应。服务端层面FastAPI的StreamingResponse本身没有超时限制但如果你用了Uvicorn它的--timeout-keep-alive参数默认是5秒这个参数控制的是连接保持时间不是流式传输时间一般不用改。真正要关注的是LLM调用本身的超时设置。location /chat/stream { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; proxy_send_timeout 300s; chunked_transfer_encoding on; }除了超时还要处理心跳。如果LLM生成过程中有较长的停顿比如检索阶段耗时较长中间层可能误判为空闲连接而断开。解决方案是定期发送注释行作为心跳。async def event_generator(query: str): yield : heartbeat\n\n async for token in rag_pipeline.astream(query): yield fdata: {json.dumps({token: token})}\n\n yield data: [DONE]\n\n以冒号开头的行是SSE的注释客户端会忽略但能保持连接活跃。我一般每15秒发一次心跳这个间隔要小于中间层的超时时间。6. 第一次闭环跑通后的验证与调优6.1 怎么判断闭环真的通了闭环跑通的标志不是接口返回了200而是从提问到答案的完整链路每个环节都有输出且合理。我一般用三个检查点来验证第一个检查点是检索结果。打印出向量检索和关键词检索各自返回的Top-5文档人工看一眼是否跟问题相关。如果检索结果明显不相关后面的生成再流畅也没用。第二个检查点是上下文组装。打印出最终拼进Prompt的上下文内容检查是否包含了回答问题的关键信息token数是否在预算内父窗口回填是否生效。第三个检查点是流式输出。用curl或者浏览器开发者工具看SSE流确认token是逐个到达的不是一次性到达的。如果是一次性到达说明中间层有缓冲。curl -N http://localhost:8000/chat/stream?query检索主链路是什么-N参数关闭curl的缓冲能实时看到SSE输出。如果输出是逐行出现的说明流式没问题。6.2 检索效果不好的常见原因闭环跑通后最常见的问题就是检索效果不好。我总结了几类原因和对应的排查方法切块粒度问题。块太大检索精度低块太小上下文碎片化。排查方法是打印命中块的原文看是否包含完整语义。如果块被切断了句子说明切块策略有问题。嵌入模型不匹配。中文场景用英文嵌入模型效果会打折扣。排查方法是拿几个典型query看Top-1的相似度分数。如果普遍低于0.7可能是模型不匹配。查询改写缺失。用户提问往往很口语化直接拿去做向量检索效果不好。加一个查询改写节点用LLM把口语化问题改写成检索友好的形式能明显提升召回率。元数据过滤缺失。如果知识库有分类、时间等元数据检索时应该加过滤条件。比如用户问最新的政策应该按时间倒序过滤而不是纯向量检索。问题现象可能原因排查方法解决方案召回文档不相关嵌入模型不匹配看Top-1相似度分数换中文嵌入模型精确术语召回差缺关键词检索检查BM25是否生效加自定义词典上下文碎片化切块太小看命中块原文父窗口回填答案缺关键信息预算截断看上下文token数调大预算或压缩流式中断中间层超时看Nginx日志调大超时心跳6.3 延迟优化的几个着力点检索主链路的延迟主要来自四块向量检索、关键词检索、重排序、LLM生成。前三个是检索阶段通常占总延迟的20%到30%LLM生成占大头。向量检索的优化用HNSW索引代替IVF查询速度能快3到5倍代价是内存占用增加。如果向量库支持GPU加速开启GPU能再快一个数量级。关键词检索的优化Elasticsearch的BM25查询本身很快瓶颈通常在分词。把分词结果缓存起来同样的query不用重复分词。重排序的优化交叉编码重排是延迟大户。如果延迟敏感可以只对Top-10做重排或者用更小的重排模型。我实测bge-reranker-base在CPU上跑Top-20要300毫秒换成bge-reranker-small能降到100毫秒以内精度损失约3%。LLM生成的优化用流式输出本身就是一种延迟优化用户感知到的首token时间比完整响应时间重要得多。另外Prompt尽量精简系统提示词不要写太长能减少首token延迟。提示我习惯在检索主链路的每个节点加耗时打点用Python的time.perf_counter()记录每个节点的开始和结束时间最后汇总打印。这样一眼就能看出瓶颈在哪个环节。生产环境可以用OpenTelemetry做分布式追踪但开发阶段打点足够了。7. 几个容易忽略的工程细节7.1 空检索结果的兜底话术检索返回空结果时不要让LLM硬编。我见过有的实现检索为空还是把空上下文拼进Prompt结果LLM开始胡编。正确的做法是检索为空时直接返回预设的兜底话术比如抱歉我没有找到相关信息您可以换个说法再问一次。如果一定要让LLM参与也要在Prompt里明确说明如果没有找到相关信息请直接告知用户不要编造。但实测下来直接返回兜底话术更可控。7.2 多轮对话的状态传递第一次闭环通常只处理单轮问答。如果要支持多轮需要把历史对话也纳入状态管理。我的做法是在状态里加一个history字段存最近N轮的问答对。检索时把当前query和上一轮query拼接起来做检索能提升指代消解的效果。但history不能无限增长否则会挤占上下文预算。我一般只保留最近3轮超过的截断。另外history里的内容也要计入token预算。7.3 日志与可观测性检索主链路的日志要记录query原文、改写后的query、向量检索Top-5的文档ID和分数、关键词检索Top-5的文档ID和分数、融合后的Top-5、最终上下文的token数、LLM的首token延迟和总延迟。这些数据是后续调优的依据。我用的日志格式是JSON Lines每行一个JSON对象方便用jq或者ELK分析。关键字段加索引方便按query或者文档ID检索。import logging import json logger logging.getLogger(rag_pipeline) def log_retrieval(query, vector_docs, keyword_docs, merged_docs): logger.info(json.dumps({ event: retrieval, query: query, vector_top5: [d[id] for d in vector_docs[:5]], keyword_top5: [d[id] for d in keyword_docs[:5]], merged_top5: [d[id] for d in merged_docs[:5]], vector_scores: [d[score] for d in vector_docs[:5]], keyword_scores: [d[score] for d in keyword_docs[:5]] }, ensure_asciiFalse))这套日志跑一段时间后你就能看出哪些query的检索效果差针对性地优化。比如发现某类query的向量分数普遍偏低可能是嵌入模型对这类语义不敏感考虑换模型或者加查询改写。7.4 配置外置与热更新检索主链路的参数很多Top-K、融合权重、预算大小、超时时间、重排开关。这些参数不要硬编码在代码里要外置到配置文件或者配置中心。我用的YAML配置文件启动时加载支持通过接口热更新。热更新这个功能在调参阶段特别有用。改一个权重不用重启服务调完立即生效效率高很多。但要注意热更新时的并发安全用读写锁保护配置对象。retrieval: vector_top_k: 20 keyword_top_k: 20 fusion_weights: vector: 0.7 keyword: 0.3 rerank: enabled: true top_n: 10 context: max_tokens: 8000 reserve_output: 2000 sse: heartbeat_interval: 15 timeout: 300这份配置是我项目里的实际配置你可以根据自己的场景调整。max_tokens设8000是因为我用的模型上下文窗口是16K留一半给输出和系统提示词。如果你的模型窗口更大可以相应调大。8. 从闭环到可用的最后一步第一次闭环跑通意味着你的RAG系统从能检索进化到了能问答。但这离好用还有距离。我建议在闭环跑通后立刻做一轮端到端测试准备20到30个典型问题覆盖事实查询、对比查询、多跳查询、否定查询等类型人工评估每个问题的答案质量。评估维度包括检索是否召回了正确文档、答案是否准确、是否有编造、流式输出是否流畅。把不达标的case记录下来逐个分析原因。大部分问题都能归到前面讲的几类切块、嵌入模型、查询改写、预算、重排。我自己的经验是第一轮端到端测试通常有30%到40%的case不达标。经过两到三轮调优能降到10%以内。剩下的10%往往是知识库本身缺失的内容或者问题本身有歧义这类case不用强求。最后分享一个我踩过的坑不要在检索主链路里做太多事。我一开始把查询改写、多路召回、重排、压缩全塞在一个图里结果调试极其困难一个环节出问题整条链路都挂。后来拆成了两个图检索图负责召回和排序生成图负责上下文组装和LLM调用。两个图通过状态传递衔接调试时可以单独跑检索图看召回效果单独跑生成图看生成质量。这个拆分让排查效率提升了很多。检索主链路是整个RAG系统的骨架骨架搭好了后面加缓存、加多轮、加权限控制都是在这个骨架上挂东西。第一次闭环不用追求完美先让它跑起来再逐步优化。跑通比完美重要。
返回列表