三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

LangGraph状态机与多源异构RAG:构建可编排的复杂智能问答系统

LangGraph状态机与多源异构RAG:构建可编排的复杂智能问答系统

1. 项目概述:当LangGraph状态机遇上多源异构RAG

最近在折腾一个智能问答系统的重构,核心目标很明确:让AI不仅能回答得准,还要能根据复杂、多步骤的用户问题,像人一样有条理地思考和行动。传统的RAG(检索增强生成)框架在处理单一知识库的简单问答时很给力,但一旦问题涉及多个数据源(比如内部文档、数据库、实时API),或者需要分步决策(比如先查政策,再算数据,最后生成报告),就显得力不从心了。这正是“系统架构设计-LangGraph状态机与多源异构RAG”这个项目要啃的硬骨头。

简单来说,这个架构是把两个强大的概念拧在了一起。一边是多源异构RAG,它解决了“信息从哪里来”的问题。我们不再依赖单一的向量数据库,而是构建一个能同时对接PDF文档、结构化数据库表、网页内容甚至实时接口的检索层,确保回答有最全面、最新鲜的依据。另一边是LangGraph状态机,它解决了“任务怎么执行”的问题。LangGraph允许我们用图(Graph)的方式定义AI的工作流,每个节点是一个处理步骤(如检索、推理、调用工具),节点间的连线定义了状态流转的逻辑。这就把一个复杂的AI任务,变成了一个可视、可控、可调试的流程。

这套组合拳特别适合需要强逻辑、多步骤交互的场景。比如,一个金融客服机器人,用户问“帮我对比一下A基金和B基金最近一年的表现,并给出风险评估”。这个任务就需要:1. 从基金文档库检索产品说明书(RAG),2. 调用实时行情API获取净值数据(工具调用),3. 从研报库获取历史风险分析(多源RAG),4. 综合所有信息生成对比报告(状态机协调)。如果你正在构建类似的需要串联多个AI能力或数据源的智能体(Agent)、复杂问答系统或自动化流程,这个架构会给你带来全新的思路和实实在在的效率提升。

2. 核心架构设计思路与选型考量

2.1 为什么是“状态机”而不是“链式调用”?

在LangChain等早期框架中,我们习惯用“链”(Chain)来组织任务。链是线性的,一个接一个执行,虽然简单,但缺乏灵活性。当任务出现分支(比如根据检索结果决定下一步是查数据库还是直接生成答案)、循环(比如信息不足时需要反问用户)或并行处理时,链就显得非常笨拙。

状态机(State Machine)模型完美适配了这种复杂性。在状态机中,我们把整个系统看作一系列“状态”(State)的集合,以及触发状态间转换的“条件”(Condition)。在LangGraph的语境下,每个“状态”对应图中的一个节点(Node),它执行特定的函数;状态之间的“边”(Edge)则定义了基于当前执行结果,下一步应该走到哪个节点。这带来了几个关键优势:

  1. 显式控制流:整个工作流的路径一目了然,不再是黑盒。你可以清晰地看到,在“检索”节点之后,系统会根据检索结果的质量(如相关度分数),决定是进入“精炼查询”节点重新检索,还是进入“生成”节点准备回答。
  2. 内置循环与条件分支:LangGraph原生支持循环(通过将边指向之前的节点)和条件边(conditional edges),这使得实现“直到找到满意答案为止”或“如果A则做B,否则做C”的逻辑变得异常简单。
  3. 状态持久化与共享:整个工作流有一个共享的“状态”(State)对象,它随着流程推进在各个节点间传递和修改。这意味着“检索”节点找到的文档,可以轻松地被后续的“分析”节点使用,无需复杂的参数传递机制。

选择LangGraph来实现这个状态机,是因为它深度集成了LangChain生态,对AI原生的工作流(调用LLM、使用工具)支持得最好。它的StateGraphCompiledGraph抽象,让定义和运行一个复杂的图变得像搭积木一样直观。

2.2 多源异构RAG的挑战与分层设计

“多源异构”听起来高大上,其实痛点很具体:你的数据散落在各处,格式还都不一样。可能有一部分是公司内部的Word/PDF文档(非结构化文本),一部分是MySQL或PostgreSQL里的业务数据(结构化数据),还有一部分需要从Confluence或某个内部API实时获取(半结构化或实时数据)。传统的单一向量检索面对这种局面会失效,因为你无法把所有东西都塞进一个向量数据库,即使能塞,数据同步和更新也是噩梦。

因此,我们的架构必须采用分层检索统一调度的策略。核心设计如下:

  1. 源适配层:为每种数据源开发一个专用的“连接器”(Connector)或“检索器”(Retriever)。例如:
    • 文档检索器:处理PDF、Word、TXT。核心是文本提取、分块(Chunking)和嵌入(Embedding),存入向量数据库(如Chroma、Weaviate)。
    • 数据库检索器:连接业务数据库。对于自然语言查询,需要将其转换为SQL(Text-to-SQL),执行查询后,将结果转换为自然语言描述。
    • API检索器:封装对内部或第三方API的调用。根据用户查询,构造API请求参数,获取JSON/XML响应并解析。
  2. 路由与调度层:这是大脑。它接收用户问题,并决定应该去查询哪个或哪几个数据源。这里可以设计一个轻量级的“路由分类器”,通常是一个经过微调的轻量级LLM或一个基于规则的决策树。例如,问题中包含“客户”、“订单”、“销售额”等关键词,则优先路由到数据库检索器;问题关于“产品手册”、“操作指南”,则路由到文档检索器。
  3. 结果融合层:当从多个源检索到结果后,需要将它们融合成一个连贯、去重、排序的上下文。这里涉及的技术包括:
    • 重排序(Re-ranking):使用专门的交叉编码器模型(如bge-reranker)对初步检索到的所有片段进行相关性重排,确保最相关的信息排在最前面。
    • 结果合成:将来自不同源的文本、数据表格、API返回的字段,整合成一段LLM容易理解的提示词(Prompt)上下文。

注意:多源RAG不是简单地把所有检索器并行跑一遍然后合并结果。不加选择地检索所有源,会导致上下文窗口被大量无关信息挤占,增加成本并降低答案质量。智能路由是保证效率和效果的关键。

2.3 LangGraph与多源RAG的融合点

那么,LangGraph状态机如何与这个多源RAG架构结合呢?答案是:将每一个关键的RAG环节,都设计为LangGraph图中的一个节点

我们可以设计一个核心工作流,其状态(State)包含诸如user_querycurrent_retrieval_resultsdecided_sourcefinal_answer等字段。图的流程可能如下:

  1. 节点:问题分析与路由。接收用户原始问题,调用一个轻量级LLM或规则引擎进行分析,输出一个路由决策(例如:{"source": ["database", "manual"], "intent": "query_sales_data"})。这个决策会被写入状态。
  2. 条件边:根据路由决策,图会走不同的分支。如果决策中包含“database”,则流向“数据库检索”节点;如果包含“manual”,则流向“文档检索”节点。这些分支可以并行执行。
  3. 节点:多源检索执行。这里可能并行存在“数据库检索节点”和“文档检索节点”。它们读取状态中的路由决策和用户问题,调用对应的检索器,将检索结果写回状态(如state[“db_results”],state[“doc_results”])。
  4. 节点:结果融合与重排序。等待并行检索节点完成后,此节点被触发。它读取状态中来自不同源的结果,执行重排序和合成,生成一个高质量的combined_context
  5. 节点:生成最终答案。将融合后的上下文和用户问题一起,发送给大语言模型(LLM),生成最终答案,并写入state[“final_answer”]

通过这样的设计,我们得到了一个可编排、可观测、可维护的智能系统。每个节点的输入输出清晰,整个决策流程像流程图一样可视,无论是调试错误还是增加新的数据源(只需增加一个节点和相应的路由逻辑),都变得非常容易。

3. 核心模块拆解与实操要点

3.1 定义LangGraph的状态(State)

状态是整个工作流运行时共享的内存,它定义了图中流动的数据结构。在LangGraph中,我们通常使用TypedDict或Pydantic模型来定义State。

from typing import TypedDict, List, Optional, Annotated from langgraph.graph.message import add_messages import operator class GraphState(TypedDict): # 用户输入 user_query: str # 路由决策 route_decision: Optional[dict] # 各源检索结果 db_retrieval_result: Optional[List[str]] doc_retrieval_result: Optional[List[str]] api_retrieval_result: Optional[dict] # 融合后的上下文 fused_context: Optional[str] # 最终答案 final_answer: Optional[str] # 用于记录对话历史(如果需要多轮) messages: Annotated[list, add_messages]

这里的关键是Annotated[list, add_messages],这是LangGraph提供的一个特殊注解,用于自动管理对话历史列表,非常方便。其他字段则根据我们的业务需求自定义。一个设计良好的State是成功的一半,它需要涵盖工作流中所有节点可能产生和消费的数据。

3.2 构建多源检索器(Retriever)

这是RAG的核心。我们以文档检索器和数据库检索器为例。

文档检索器(基于向量数据库)

from langchain_community.vectorstores import Chroma from langchain_openai import OpenAIEmbeddings from langchain.text_splitter import RecursiveCharacterTextSplitter from langchain_community.document_loaders import PyPDFLoader class DocumentRetriever: def __init__(self, persist_directory="./chroma_db"): self.embeddings = OpenAIEmbeddings(model="text-embedding-3-small") self.vectorstore = Chroma( persist_directory=persist_directory, embedding_function=self.embeddings ) self.text_splitter = RecursiveCharacterTextSplitter( chunk_size=1000, chunk_overlap=200 ) def add_documents(self, file_path): """向知识库添加新文档""" loader = PyPDFLoader(file_path) documents = loader.load() splits = self.text_splitter.split_documents(documents) self.vectorstore.add_documents(splits) def retrieve(self, query: str, k: int = 4) -> List[str]: """检索相关文档片段""" docs = self.vectorstore.similarity_search(query, k=k) return [doc.page_content for doc in docs]

数据库检索器(基于Text-to-SQL)

from langchain_community.utilities import SQLDatabase from langchain.chains import create_sql_query_chain from langchain_openai import ChatOpenAI class DatabaseRetriever: def __init__(self, db_uri): self.db = SQLDatabase.from_uri(db_uri) self.llm = ChatOpenAI(model="gpt-4", temperature=0) self.query_chain = create_sql_query_chain(self.llm, self.db) def retrieve(self, natural_language_query: str) -> str: """将自然语言转换为SQL并执行,返回自然语言结果描述""" # 1. 生成SQL sql_query = self.query_chain.invoke({"question": natural_language_query}) # 2. 执行SQL(这里需要谨慎,最好有权限控制和验证) result = self.db.run(sql_query) # 3. 将结果转换为易于理解的文本 summary = f"根据你的问题“{natural_language_query}”,查询到的数据结果是:{result}" return summary

实操心得:对于数据库检索器,直接让LLM生成并执行SQL是高风险操作。在生产环境中,务必加入以下安全层:1) SQL语法验证;2) 只读权限数据库连接;3) 查询结果行数限制;4) 敏感表/字段过滤。更好的做法是使用“语义层”或“预定义查询模板”,将自然语言映射到安全的查询语句上。

3.3 实现智能路由节点

路由节点负责分析用户意图,决定查询路径。我们可以用一个简单的基于LLM的分类器来实现。

from langchain_core.prompts import ChatPromptTemplate from langchain_openai import ChatOpenAI class RouterNode: def __init__(self): self.llm = ChatOpenAI(model="gpt-3.5-turbo", temperature=0) self.prompt = ChatPromptTemplate.from_messages([ ("system", """你是一个智能路由助手。请分析用户问题,判断它最可能涉及哪类数据源。 可用的数据源类型有: - `document`: 公司内部文档、手册、PDF文件。 - `database`: 业务数据,如销售记录、用户信息、产品库存。 - `api`: 需要调用外部或内部API获取的实时信息,如天气、股价、物流状态。 如果问题明显涉及多个类型,可以返回多个。 请只返回一个JSON对象,格式如:{{"sources": ["document", "database"], "primary_intent": "query_product_info"}}"""), ("human", "{question}") ]) self.chain = self.prompt | self.llm def route(self, state: GraphState) -> dict: """路由决策函数,将被LangGraph节点调用""" question = state["user_query"] response = self.chain.invoke({"question": question}) # 解析LLM返回的JSON import json try: decision = json.loads(response.content) except: decision = {"sources": ["document"], "primary_intent": "general"} # 默认降级 return {"route_decision": decision}

这个节点接收State,从中取出用户问题,调用LLM进行分析,然后将结构化的路由决策写回State。在图中,下一个节点或条件边就可以读取state[“route_decision”]来决定后续流程。

4. 组装LangGraph状态机与工作流编排

4.1 构建图节点与边

有了核心组件,现在用LangGraph把它们组装起来。我们首先创建各个节点函数,然后定义它们之间的流转关系。

from langgraph.graph import StateGraph, END from langgraph.graph import MessagesState # 假设我们已经有了上述类的实例 router = RouterNode() doc_retriever = DocumentRetriever() db_retriever = DatabaseRetriever("sqlite:///./test.db") # 1. 定义节点函数 def route_question(state: GraphState): """路由节点""" return router.route(state) def retrieve_from_docs(state: GraphState): """文档检索节点""" if "document" in state.get("route_decision", {}).get("sources", []): query = state["user_query"] results = doc_retriever.retrieve(query, k=3) return {"doc_retrieval_result": results} return {"doc_retrieval_result": None} # 如果路由未指定,返回None def retrieve_from_db(state: GraphState): """数据库检索节点""" if "database" in state.get("route_decision", {}).get("sources", []): query = state["user_query"] results = db_retriever.retrieve(query) return {"db_retrieval_result": results} return {"db_retrieval_result": None} def fuse_results(state: GraphState): """结果融合节点""" contexts = [] if state.get("doc_retrieval_result"): contexts.append("[来自文档的知识]:\n" + "\n---\n".join(state["doc_retrieval_result"])) if state.get("db_retrieval_result"): contexts.append("[来自数据库的数据]:\n" + state["db_retrieval_result"]) if not contexts: fused = "未检索到相关信息。" else: fused = "\n\n".join(contexts) return {"fused_context": fused} def generate_answer(state: GraphState): """答案生成节点""" from langchain_core.prompts import ChatPromptTemplate from langchain_openai import ChatOpenAI llm = ChatOpenAI(model="gpt-4", temperature=0.3) prompt = ChatPromptTemplate.from_messages([ ("system", "你是一个专业的助手,请根据以下提供的上下文信息,准确、有条理地回答用户的问题。如果上下文信息不足以回答问题,请如实说明。"), ("human", "上下文信息:\n{context}\n\n用户问题:{question}") ]) chain = prompt | llm response = chain.invoke({ "context": state.get("fused_context", "无相关信息"), "question": state["user_query"] }) return {"final_answer": response.content} # 2. 创建状态图 workflow = StateGraph(GraphState) # 3. 添加节点 workflow.add_node("router", route_question) workflow.add_node("retrieve_docs", retrieve_from_docs) workflow.add_node("retrieve_db", retrieve_from_db) workflow.add_node("fuser", fuse_results) workflow.add_node("generator", generate_answer) # 4. 设置入口和边 workflow.set_entry_point("router") # 从router出来后,并行执行文档和数据库检索 workflow.add_edge("router", "retrieve_docs") workflow.add_edge("router", "retrieve_db") # 设置一个条件,等待两个检索节点都完成(或跳过)后,再进入融合节点 # LangGraph 0.2+ 版本提供了更优雅的并行和汇聚控制,这里我们用条件边模拟 def after_retrieve(state): # 简单的逻辑:只要路由决策存在,就进入融合节点。 # 更复杂的逻辑可以检查各个检索结果是否已完成。 if state.get("route_decision"): return "fuser" return END workflow.add_conditional_edges( "retrieve_docs", after_retrieve, {"fuser": "fuser", END: END} ) # 同样为retrieve_db添加边指向fuser(实际中需要更精细的汇聚逻辑) workflow.add_edge("retrieve_db", "fuser") # 从融合节点到生成节点,再到结束 workflow.add_edge("fuser", "generator") workflow.add_edge("generator", END) # 5. 编译图 app = workflow.compile()

4.2 运行与调试工作流

图编译好后,就可以像调用函数一样运行它。传入初始状态,获取最终结果。

# 准备初始状态 initial_state = GraphState( user_query="上一季度我们产品A在华北区的销售额是多少?另外,产品A的用户手册里提到的最大负载是多少?", messages=[] # 初始化对话历史 ) # 运行图 final_state = app.invoke(initial_state) print("最终答案:", final_state["final_answer"]) print("\n--- 完整状态追踪 ---") for key, value in final_state.items(): if key != "messages": # 过滤掉可能很长的历史 print(f"{key}: {value}")

LangGraph的一个强大功能是可视化。你可以将图导出为PNG,直观地看到整个工作流。

# 导出图结构(需要安装graphviz) from IPython.display import Image, display try: display(Image(app.get_graph().draw_mermaid_png())) except: # 或者打印文本表示 print(app.get_graph().draw_ascii())

5. 高级特性与生产级考量

5.1 实现子图(Subgraph)进行模块化

当工作流变得非常复杂时,可以将一部分功能封装成子图。例如,整个“多源检索与融合”可以作为一个子图,在主图中只用一个节点表示。这极大地提升了可维护性和复用性。

from langgraph.graph import StateGraph # 创建一个“检索融合”子图 retrieval_subgraph_builder = StateGraph(GraphState) retrieval_subgraph_builder.add_node(“retrieve_docs”, retrieve_from_docs) retrieval_subgraph_builder.add_node(“retrieve_db”, retrieve_from_db) retrieval_subgraph_builder.add_node(“fuser”, fuse_results) # ... 设置子图内部的边 retrieval_subgraph = retrieval_subgraph_builder.compile() # 在主图中,将子图作为一个节点添加 workflow.add_node(“retrieval_fusion_subgraph”, retrieval_subgraph)

5.2 引入人工干预与持久化检查点

对于关键业务,有时需要“人在环路”(Human-in-the-loop)。LangGraph支持持久化检查点(Persisted Checkpoints),允许你将工作流在任何节点的状态保存下来。例如,你可以在“生成答案”节点前设置一个检查点,将检索到的上下文发送给人工审核,审核通过后再继续执行生成。

from langgraph.checkpoint import MemorySaver # 在编译图时加入检查点存储器 memory = MemorySaver() app = workflow.compile(checkpointer=memory) # 运行到某个节点后暂停,获取一个线程ID和检查点 config = {"configurable": {"thread_id": "user_123_session_1"}} initial_state = GraphState(user_query="...", messages=[]) # 假设我们只想运行到‘fuser’节点前 # 可以通过自定义边或中断逻辑实现,这里展示概念 # 保存的状态可以后续被加载和继续执行

5.3 性能优化与监控

在生产环境中,以下几点至关重要:

  1. 检索优化

    • 索引策略:针对文档,选择合适的文本分块大小和重叠度。太小丢失上下文,太大降低精度。
    • 混合搜索:结合向量搜索(语义相似)和关键词搜索(BM25),提升召回率。
    • 缓存:对常见的查询结果进行缓存,尤其是API和数据库查询结果,可以大幅降低延迟和成本。
  2. 图执行优化

    • 并行化:确保独立的节点(如retrieve_docsretrieve_db)真正并行执行,而不是顺序执行。LangGraph的异步支持可以帮上忙。
    • 超时与重试:为每个节点设置超时和重试机制,特别是调用外部API或复杂数据库查询的节点。
  3. 可观测性

    • 日志记录:在每个节点的入口和出口记录详细的日志,包括输入、输出、耗时和可能的错误。
    • 链路追踪:集成像OpenTelemetry这样的追踪工具,为每个用户查询生成一个完整的执行链路图,便于定位性能瓶颈和错误。

6. 常见问题排查与实战技巧

在实际开发和运维中,你会遇到各种各样的问题。下面是一个快速排查指南:

问题现象可能原因排查步骤与解决方案
答案与文档内容无关(幻觉)1. 检索到的上下文不相关。
2. LLM忽略了上下文。
1.检查检索结果:打印出fused_context,看是否与问题相关。调整检索器的k值(返回数量)或尝试重排序。
2.强化Prompt:在系统指令中明确强调“必须且仅能依据提供的上下文回答”,使用类似“If the context doesn‘t contain the answer, say ‘I don’t know‘.”的指令。
检索速度慢1. 向量数据库索引未优化。
2. 并行检索未生效。
3. 文本分块或嵌入模型太慢。
1. 确保向量数据库使用了合适的索引(如HNSW)。
2. 检查LangGraph节点配置,确保retrieve_docsretrieve_db是并行边,而非顺序边。
3. 考虑使用更快的嵌入模型(如text-embedding-3-small),或对文档进行预处理和预嵌入。
数据库检索返回错误或空结果1. Text-to-SQL生成的SQL语法错误。
2. 自然语言问题歧义导致查询不准。
1.增加SQL验证层:在执行前用sqlparse等库检查SQL语法,或使用EXPLAIN预览。
2.提供数据库Schema提示:在给LLM的Prompt中,加入精简的、相关的表结构描述,大幅提升SQL生成准确率。
LangGraph图编译或运行报错1. 状态(State)字段定义与节点返回值不匹配。
2. 边(Edge)指向不存在的节点。
1.仔细核对State的TypedDict定义,确保每个节点返回的字典键名都能在State中找到对应。
2.使用app.get_graph().draw_ascii()可视化图,检查所有边的指向是否正确。从简单的两个节点开始测试,逐步增加复杂度。
多轮对话中状态混乱1. 对话历史(messages)未正确更新。
2. 前一轮的检索结果污染了当前轮状态。
1. 利用好Annotated[list, add_messages],它会在调用LLM的节点自动管理历史。对于不涉及LLM的节点,注意不要误修改messages字段。
2. 在每一轮对话开始时,考虑有选择地重置State。例如,可以设计一个reset_state节点,在收到新问题时,只保留messages历史,清空retrieval_result等中间字段。

最后分享一个实战技巧:从简单开始,迭代演进。不要试图第一次就设计出包含所有数据源和复杂分支的完美状态机。我的建议是:

  1. 第一步:先实现一个单源的、线性的RAG流程(用户问题 -> 检索 -> 生成),让它跑通。
  2. 第二步:引入LangGraph,把这个线性流程改造成一个最简单的两节点状态机(检索节点 -> 生成节点),感受状态流转。
  3. 第三步:增加第二个数据源,并加入路由节点,实现条件分支。
  4. 第四步:引入并行、子图、检查点等高级特性。

这种渐进式的方法,能让复杂系统的构建过程更可控,也更容易调试。毕竟,一个清晰、健壮的状态机图,其价值不仅在于让AI更聪明,也在于让开发和维护它的我们,思路更清晰。

← 返回列表