LangGraph状态机:构建企业级多源异构RAG系统的工程实践

发布时间:2026/8/8 11:11:26
LangGraph状态机:构建企业级多源异构RAG系统的工程实践 1. 项目缘起当RAG遇上复杂业务流最近在做一个企业级知识问答系统的重构需求听起来挺常规用户提问系统从文档库里找答案。但一深入坑就来了。文档来源五花八门——有躺在Confluence里的产品手册有GitHub上的技术文档有飞书里的会议纪要甚至还有一堆历史遗留的PDF和Word文件。这还不算完用户的问题也千奇百怪有的简单到直接检索就能回答有的则需要拆解成多个子问题分别查询不同来源再把结果汇总、去重、排序最后生成一个连贯的答案。最开始我们用一个简单的LangChain链条Chain硬扛代码很快就变成了“意大利面条”——一堆if-else嵌套状态管理混乱错误处理更是灾难。每当要新增一个数据源或者调整查询逻辑都感觉是在拆弹。这时候我意识到我们需要的不是一个更复杂的链条而是一个能清晰描述和控制“对话状态”与“执行流程”的架构。这让我把目光投向了LangGraph特别是其状态机StateGraph模型用它来驾驭多源异构RAG检索增强生成这个“怪兽”成了我们架构升级的核心。简单说这个项目就是用LangGraph的状态机为混乱的多源检索流程建立秩序。它不再是一条道走到黑而是一个有明确节点如“解析问题”、“检索A源”、“检索B源”、“综合判断”和流转规则边的图。系统当前“在哪儿”、“知道什么”、“下一步该干嘛”都清晰定义在一个状态对象里。这样一来无论是扩展新数据源还是实现复杂的多步推理Agentic RAG都变得模块化且可控。下面我就把这套架构的设计思路、核心实现以及踩过的坑毫无保留地分享出来。2. 核心架构LangGraph状态机驱动的工作流引擎传统的LangChain Chain或Agent其控制流往往是隐式的、线性的或者通过有限的工具调用来实现分支。在处理多源异构RAG这种带有明显阶段性和条件分支的场景时就显得力不从心。LangGraph的StateGraph则提供了显式的、基于图的状态机抽象完美匹配我们的需求。2.1 状态State设计系统的“记忆中枢”一切始于状态定义。在LangGraph中状态是一个字典或Pydantic模型它随着工作流的执行而演变。我们的核心状态设计如下from typing import TypedDict, List, Optional, Annotated from langgraph.graph.message import add_messages import operator class GraphState(TypedDict): 定义工作流的状态结构。 # 用户输入 question: str # 问题解析后的结果如意图分类、关键词、子问题列表 parsed_query: Optional[dict] # 各数据源的检索结果 retrieval_results: Annotated[dict, operator.add] # 关键这是一个可追加的字典 # 当前已尝试的检索源列表用于避免重复检索 attempted_sources: Annotated[List[str], operator.add] # 最终用于生成答案的合成上下文 final_context: Optional[str] # 最终答案 final_answer: Optional[str] # 错误信息或特殊标志 error: Optional[str] # 决定下一个节点的控制信号 next: Optional[str]设计要点解析Annotated与operator.add这是LangGraph状态更新的精髓。retrieval_results和attempted_sources被标记为Annotated[..., operator.add]。这意味着当不同节点函数返回的新状态中包含这些字段时LangGraph不会直接覆盖旧值而是会执行add操作。对于字典add相当于dict.update()对于列表相当于list.extend()。这保证了来自不同数据源的检索结果能自然地汇聚到一起而不是互相覆盖。parsed_query这是一个中间解析结果。我们可能用一个专门的LLM调用或规则引擎将原始问题解析为结构化的信息例如{intent: comparison, entities: [产品A, 产品B], sub_questions: [产品A的特性, 产品B的特性]}。这为后续的条件路由提供了依据。next字段这是一个可选的控制字段。任何节点都可以通过设置state[“next”] “node_name”来显式指定下一个要执行的节点从而实现复杂的分支跳转。如果不设置则默认按照图中定义的边来流转。2.2 节点Nodes模块化的功能单元节点就是普通的Python函数或可调用对象它接收当前状态执行一些操作然后返回更新后的状态。每个节点职责单一。示例1查询解析节点def parse_query(state: GraphState) - GraphState: 解析用户问题判断意图并提取关键信息。 question state[“question”] # 这里可以调用一个轻量级LLM如GPT-3.5-turbo或使用规则 # 假设我们调用一个LLM进行解析 from langchain_core.prompts import ChatPromptTemplate from langchain_openai import ChatOpenAI prompt ChatPromptTemplate.from_messages([ (“system”, “你是一个查询解析助手。请将用户问题解析为JSON格式包含’intent’意图、’keywords’关键词列表、’requires_multisource’是否需要多源布尔值字段。”), (“human”, “{question}”) ]) llm ChatOpenAI(model“gpt-3.5-turbo”, temperature0) parser_chain prompt | llm # 实际应用中这里需要更健壮的JSON解析和错误处理 try: parsed parser_chain.invoke({“question”: question}) # 假设LLM返回了格式正确的JSON字符串 import json parsed_dict json.loads(parsed.content) except Exception as e: parsed_dict {“intent”: “general”, “keywords”: [], “requires_multisource”: False, “error”: str(e)} return {“parsed_query”: parsed_dict}示例2特定数据源检索节点以Confluence为例def retrieve_from_confluence(state: GraphState) - GraphState: 从Confluence知识库检索相关信息。 if “confluence” in state.get(“attempted_sources”, []): # 已检索过跳过 return {“retrieval_results”: {“confluence”: “[已跳过重复检索]”}} parsed state.get(“parsed_query”) keywords parsed.get(“keywords”, []) if parsed else [] query “ “.join(keywords) if keywords else state[“question”] # 这里是具体的检索逻辑例如使用Confluence REST API或已构建的向量库 # 假设我们有一个已初始化的Confluence检索器 confluence_retriever try: docs confluence_retriever.invoke(query) # 返回List[Document] relevant_text “\n\n”.join([doc.page_content for doc in docs[:3]]) # 取Top3 result {“confluence”: relevant_text} except Exception as e: result {“confluence”: f“检索失败: {str(e)}”} # 更新状态记录结果并标记该源已尝试 return { “retrieval_results”: result, “attempted_sources”: [“confluence”] # operator.add 会将其追加到列表 }注意每个检索节点都应该检查attempted_sources避免在循环或重试中重复检索同一源浪费资源。2.3 边Edges与路由定义工作流逻辑边决定了状态在节点间的流转方向。LangGraph支持条件边Conditional Edge这是实现智能路由的关键。from langgraph.graph import StateGraph, END from langgraph.graph import START # 初始化图 workflow StateGraph(GraphState) # 1. 添加节点 workflow.add_node(“parse”, parse_query) workflow.add_node(“retrieve_confluence”, retrieve_from_confluence) workflow.add_node(“retrieve_github”, retrieve_from_github) # 假设有另一个节点 workflow.add_node(“retrieve_internal_wiki”, retrieve_from_internal_wiki) workflow.add_node(“synthesize”, synthesize_answer) # 合成答案节点 workflow.add_node(“handle_error”, handle_error) # 错误处理节点 # 2. 设置入口 workflow.set_entry_point(“parse”) # 3. 添加普通边解析后默认开始并发检索这里简化实际可能并发 workflow.add_edge(“parse”, “retrieve_confluence”) # 4. 添加条件边根据解析结果决定下一步是检索其他源还是直接合成 def route_after_retrieve(state: GraphState) - str: 决定在检索一个源之后做什么。 parsed state.get(“parsed_query”, {}) attempted state.get(“attempted_sources”, []) results state.get(“retrieval_results”, {}) # 条件1: 如果解析要求多源且还有未尝试的源则继续检索下一个源 # 这里需要一个预定义的源列表和顺序逻辑为简化假设有判断函数should_retry_next_source if parsed.get(“requires_multisource”) and should_retry_next_source(attempted): # 返回下一个检索节点的名称例如基于某种优先级 return get_next_source_node(attempted) # 条件2: 如果已有足够结果或不需要多源则进入合成阶段 elif is_results_sufficient(results) or not parsed.get(“requires_multisource”): return “synthesize” # 条件3: 其他情况如出错进入错误处理 else: return “handle_error” # 将条件边添加到某个检索节点之后例如Confluence检索后 workflow.add_conditional_edges( “retrieve_confluence”, route_after_retrieve, { “retrieve_github”: “retrieve_github”, “retrieve_internal_wiki”: “retrieve_internal_wiki”, “synthesize”: “synthesize”, “handle_error”: “handle_error”, } ) # 5. 为其他检索节点添加类似的条件边或直接边 workflow.add_edge(“synthesize”, END) workflow.add_edge(“handle_error”, END) # 6. 编译图 app workflow.compile()条件路由的精髓route_after_retrieve函数是工作流的大脑。它检查当前状态如解析意图、已尝试的源、已有结果的质量并返回下一个节点的名字。这使得工作流不再是固定的流水线而是一个能根据上下文动态调整的智能体Agentic RAG的雏形。3. 多源异构检索的工程化实践有了状态机框架接下来要解决的就是“多源异构”的具体问题。每个数据源都有其独特的访问方式、数据格式和相关性判断标准。3.1 数据源抽象与统一接口我们不能为每个源写死不同的调用代码。一个好的做法是定义一个抽象的检索器接口。from abc import ABC, abstractmethod from typing import List from langchain_core.documents import Document class BaseRetriever(ABC): 所有数据源检索器的基类。 source_name: str abstractmethod def retrieve(self, query: str, top_k: int 5) - List[Document]: 检索接口返回Document列表。 pass def format_results(self, docs: List[Document]) - str: 将Document列表格式化为字符串上下文。 return “\n\n”.join([f“【来源{self.source_name}】\n{doc.page_content}” for doc in docs]) # 具体实现Confluence检索器 class ConfluenceRetriever(BaseRetriever): source_name “confluence” def __init__(self, space_key, api_token): # 初始化Confluence客户端等 self.client ConfluenceClient(space_key, api_token) # 可能还有一个本地向量库索引 self.vectorstore Chroma(persist_directory“./confluence_index”, embedding_functionembedding_model) def retrieve(self, query: str, top_k: int 5) - List[Document]: # 策略1先用关键词在向量库中做语义检索 vector_docs self.vectorstore.similarity_search(query, ktop_k) # 策略2如果向量结果置信度低再用Confluence API的全文搜索作为后备 if self._is_low_confidence(vector_docs): api_docs self.client.search_by_cql(query, limittop_k) # 将API结果转换为Document对象 api_docs [Document(page_contentitem[‘content’], metadata{“source”: “api”}) for item in api_docs] return api_docs return vector_docs # 具体实现GitHub Wiki检索器 class GitHubWikiRetriever(BaseRetriever): source_name “github_wiki” def __init__(self, repo_owner, repo_name, access_token): self.github Github(access_token) self.repo self.github.get_repo(f“{repo_owner}/{repo_name}”) def retrieve(self, query: str, top_k: int 5) - List[Document]: # 获取wiki页面列表对页面内容进行本地向量检索需预先爬取和索引 # 或者使用GitHub的搜索API如果wiki内容公开 # ...这样在LangGraph的节点函数中我们可以通过一个统一的工厂或配置来获取对应的检索器实例使节点代码保持简洁。3.2 检索结果融合与重排序Reranking当从多个源获取到文档列表后直接拼接可能效果很差。我们需要融合和重排序。策略1按源优先级简单拼接在synthesize节点中我们可以定义一个源优先级列表例如[“internal_wiki”, “confluence”, “github”]。然后按此顺序从state[“retrieval_results”]中取出各源的格式化文本进行拼接。这种方法简单但忽略了跨源文档的相关性。策略2使用交叉编码器进行重排序这是更高级的做法。将所有检索到的文档无论来自哪个源混合成一个大的列表然后使用一个重排序模型如bge-reranker、Cohere rerank根据原始问题对所有文档进行相关性打分并重新排序。def rerank_documents(query: str, all_docs: List[Document], top_n: int 10) - List[Document]: 使用重排序模型对文档进行精排。 from FlagEmbedding import FlagReranker reranker FlagReranker(‘BAAI/bge-reranker-large’, use_fp16True) # 示例模型 pairs [(query, doc.page_content) for doc in all_docs] scores reranker.compute_score(pairs) # 得到相关性分数列表 # 将分数与文档关联并排序 scored_docs list(zip(scores, all_docs)) scored_docs.sort(keylambda x: x[0], reverseTrue) reranked_docs [doc for _, doc in scored_docs[:top_n]] return reranked_docs在synthesize节点中可以先收集所有源的Document对象调用rerank_documents再将排名靠前的文档内容格式化为final_context供LLM生成答案。这能显著提升最终答案的质量尤其是当不同源返回了相似或互补信息时。3.3 处理“零结果”与降级策略多源检索中某个源返回零结果是常事。我们的状态机需要能优雅处理。节点级处理在每个检索节点内部如果检索失败或返回空可以返回一个特定的标记结果如{“confluence”: “【未找到相关信息】”}而不是抛出异常中断流程。路由决策考虑在条件路由函数route_after_retrieve中可以检查state[“retrieval_results”]中最新源的结果质量。如果结果为空或质量极低可以决策跳过下一个同类型源或转向一个通用的网络搜索后备节点如果允许。合成阶段处理在synthesize节点中生成最终上下文时可以过滤掉那些标记为无效或空的结果避免LLM被无效信息干扰。4. 高级模式子图Subgraph与长期记忆对于极其复杂的场景比如一个多轮对话中需要多次调用这个RAG工作流或者工作流本身某些部分如“多源检索与融合”可以被复用LangGraph的子图功能就派上用场了。4.1 将RAG工作流封装为子图我们可以把上面构建的整个多源RAG流程从parse到synthesize编译成一个独立的图然后将其作为一个“超级节点”嵌入到一个更大的父图中。# 假设我们已经定义并编译了上面的多源RAG图为 rag_graph from langgraph.graph import StateGraph as ParentStateGraph class ParentState(TypedDict): messages: Annotated[List, add_messages] # 对话历史 current_question: str rag_output: Optional[dict] # ... 其他父图状态 def call_rag_subgraph(state: ParentState): 调用RAG子图的函数。 # 从父状态中提取问题 question state[“current_question”] # 初始化子图需要的状态 rag_initial_state GraphState(questionquestion) # 运行子图 rag_final_state rag_graph.invoke(rag_initial_state) # 将子图结果带回父状态 return {“rag_output”: {“answer”: rag_final_state[“final_answer”], “context”: rag_final_state[“final_context”]}} # 构建父图例如一个对话机器人 parent_workflow ParentStateGraph(ParentState) parent_workflow.add_node(“rag_agent”, call_rag_subgraph) # ... 添加其他节点如对话管理、工具调用等 parent_workflow.add_edge(START, “rag_agent”) # ... 设置更多边 parent_app parent_workflow.compile()这样复杂的RAG逻辑被封装和隔离父图专注于更高层次的对话流和任务规划代码结构更清晰。4.2 集成长期记忆Long-term Memory在多轮对话中记住之前的交互历史至关重要。LangGraph通过与状态中Annotated[List, add_messages]的配合可以很方便地集成聊天记忆。在父图状态中定义消息列表如上例的ParentState。在调用RAG子图时注入历史我们可以将对话历史中的相关部分如前几轮的问题和答案作为上下文连同当前问题一起送给RAG子图。这可能需要修改子图的GraphState增加一个conversation_history字段并在parse_query节点中考虑历史信息来更好地解析当前问题。将RAG结果追加到历史子图返回答案后父图将问题和答案作为一对消息HumanMessage/AIMessage追加到state[“messages”]中。add_messages这个operator确保了消息列表的正确累积。这种模式使得构建一个能进行深入、多轮、基于知识库对话的Agent变得非常直观。5. 实战踩坑与性能调优理论很美好但实际落地时我们遇到了不少挑战。坑1状态对象的意外共享与污染最初我们在不同的节点函数中直接修改传入的state字典如state[“key”] new_value。这在一个图被多个线程并发调用时导致了诡异的状态污染。牢记每个节点函数都应该返回一个新的字典包含你想要更新的键值对。LangGraph会帮你合并。永远把state当作不可变的输入。坑2条件边函数的复杂度失控route_after_retrieve函数最初塞满了各种业务逻辑越来越难以维护和测试。解决方案将路由逻辑拆分成多个小的、可测试的纯函数。例如单独的函数should_retrieve_github(state)、is_answer_ready(state)等。然后在主路由函数中组合这些判断。这样逻辑更清晰也便于单元测试。坑3检索节点超时导致整个图卡死某个外部数据源API不稳定偶尔超时会阻塞整个工作流的执行。解决方案设置超时在每个检索节点的具体实现中使用asyncio.wait_for或requests的超时参数。引入异步节点LangGraph支持异步节点。将IO密集型的检索节点定义为async def并在图中并发执行可以大幅减少总耗时。实现降级在节点内捕获超时异常并返回降级结果如从缓存中获取旧数据或返回一个提示“该源暂时不可用”。坑4LLM调用成本与延迟parse_query和synthesize节点都需要调用LLM是主要的成本和时间消耗点。缓存对parse_query的结果进行缓存。相同或相似的问题直接使用缓存的结果。可以使用langchain.cache如InMemoryCache,SQLiteCache或Redis。流式输出对于synthesize节点生成的最终答案使用LangChain或LangGraph的流式输出支持让用户能尽快看到答案的开头提升体验。模型选型parse_query任务相对简单可以使用更小、更快的模型如gpt-3.5-turbo甚至本地小模型。synthesize任务对质量要求高再用大模型如gpt-4。性能调优建议可视化与调试使用app.get_graph().draw_mermaid_png()输出流程图帮助理解复杂的工作流。在开发时传入config{“recursion_limit”: 50, “configurable”: {“thread_id”: “test1”}}来运行并打印中间状态是调试的不二法门。并发执行对于彼此独立的检索节点如retrieve_confluence和retrieve_github如果它们之间没有严格的先后顺序可以利用LangGraph的Pregel引擎的并发特性。通过巧妙设计边让它们从同一个节点出发实现并行检索最后再汇聚到下一个节点进行结果融合。设置中断Interrupt在某些情况下你可能需要提前终止图的执行。虽然compiled_graph的stream()方法本身没有直接的“终止”API但你可以通过在外层包装异步任务或者在状态中设置一个should_stop标志并在条件边函数中检查这个标志来跳转到END节点实现软中断。6. 与LangChain的对比及选型思考很多人会问有了LangChain为什么还要用LangGraph我的体会是LangChain是“乐高积木”而LangGraph是“积木的组装说明书”。LangChain提供了极其丰富的组件Models, Prompts, Chains, Agents, Tools, Retrievers, Memory等让你能快速搭建起一个AI应用的原型。它的Chain和Agent对于线性或简单循环的任务足够了。LangGraph当你需要描述一个非线性的、有状态的、带条件分支的复杂工作流时它就成为了必需品。它的状态机模型将“控制流”和“数据流”显式化、可视化。这对于构建健壮的、可维护的、涉及多步骤决策比如我们这里的多源条件检索的AI应用至关重要。选型指南如果你的任务是一条直线输入 - 检索 - 生成 - 输出用LangChain Expression Language (LCEL)构建一个简单的Chain就够了。如果你的任务需要根据中间结果决定下一步做什么或者需要在多个可能路径间循环比如多工具协作的Agent、多轮审核流程、复杂的决策树式问答那么LangGraph的状态机是你的最佳选择。它带来的代码结构清晰度和流程可控性在项目复杂度上升时会拯救你。回到我们的“多源异构RAG”项目LangGraph状态机不仅帮我们理清了混乱的流程其模块化的设计也让团队协作变得容易。前端同事可以专注于设计更好的parse_query策略后端同事可以不断接入新的BaseRetriever实现而整个工作流的骨架图定义保持稳定。当产品经理提出“如果A源没结果能不能先查B源再查C源如果还是不行就转人工”这种需求时我们只需要在route_after_retrieve函数里加几个条件判断而不用重写整个应用逻辑。这种灵活性和掌控感是传统链式编程难以比拟的。