LangGraph实战:构建工业级多智能体工作流与RAG集成

📅 2026/8/18 19:25:01
LangGraph实战:构建工业级多智能体工作流与RAG集成
最近在尝试把一些单点任务升级成自动化流程时我遇到了一个典型问题用脚本或简单工具链跑通一次任务很容易但一旦想让多个步骤、多个模型、多个判断点串联起来形成一个能稳定运行、状态可控、还能处理异常的“工业级”流程立刻就变得棘手起来。脚本越写越长状态管理混乱错误处理像打补丁想加个新功能或换个模型都得大动干戈。这让我开始重新审视那些宣称能构建“智能体”的框架。很多工具确实能快速拼出一个原型演示起来很酷但离真正的“工业级”落地中间还隔着状态管理、错误恢复、多智能体协作、长期记忆和可观测性这几座大山。直到我开始深入使用 LangGraph并结合 LangChain 的生态才感觉找到了一个能把想法系统化落地的路径。它提供的不是另一个花哨的演示而是一套用于构建可控、可维护、可扩展状态化工作流的“工程骨架”。今天我们就来深度拆解这套架构。我不会只讲概念而是会结合一个从零到一的完整实战项目带你走过这几个关键阶段如何用 StateGraph 设计并管控复杂状态流、如何将单智能体扩展为分工明确的多智能体系统、如何无缝集成 RAG 与多种大模型最终形成一个接近生产可用的智能体应用。你会发现真正的价值不在于调用一两个 API而在于构建一个可靠、透明且易于迭代的自动化系统。1. 从“一次性脚本”到“状态化工作流”为什么需要 LangGraph在开始写代码之前我们必须先想清楚一个问题当我们在说“构建一个智能体”时我们到底在构建什么是写一个调用了大模型 API 的函数吗是串联几个工具调用吗这些都很初级。一个工业级的智能体本质上是一个有状态、可分支、可循环、具备错误恢复能力的工作流引擎。想象一个客服场景的智能体用户输入问题 - 判断意图分类- 如果是查询知识库则检索 - 组织答案 - 如果答案不完整可能触发二次检索或转人工 - 最终回复。这个流程里有明确的“状态”当前进行到哪一步、手里有什么数据有“分支”不同的意图走向不同的处理节点有“循环”检索-评估-再检索还需要处理各种异常检索失败、模型超时。如果你用纯脚本写很快就会变成“面条代码”状态用全局变量传来传去循环和分支靠一堆if-else硬编码加个新步骤或者改个逻辑都心惊胆战。而 LangGraph 的核心抽象——StateGraph就是为了解决这个问题而生的。1.1 StateGraph把工作流画成一张“状态转移图”StateGraph 鼓励你用“图”的思维来建模工作流。图中的节点Node是一个个执行单元比如意图判断、知识检索、答案生成边Edge定义了节点之间的流转条件。整个系统的“状态”是一个共享的数据结构在各个节点间传递和修改。这种设计带来了几个关键优势可视化与可理解性工作流不再是隐藏在代码里的逻辑而是一张可以画出来的图。新人能快速理解业务逻辑。模块化每个节点功能单一易于开发、测试和复用。修改一个节点不会轻易影响全局。显式状态管理所有中间数据都存在于一个明确定义的“状态”对象中避免了隐式依赖和全局变量的混乱。内置的控制流LangGraph 原生支持条件分支conditional_edge和循环将边指向之前的节点这让实现复杂的业务逻辑变得直观。1.2 LangGraph 与 LangChain 的关系生态与引擎很多人会混淆 LangGraph 和 LangChain。简单来说LangChain是一个丰富的“生态工具箱”提供了连接大模型、向量数据库、工具链、记忆模块等各类组件的标准化接口和实现。它让你能方便地“使用”各种AI能力。LangGraph则是建立在 LangChain 之上的“工作流引擎”或“编排框架”。它利用 LangChain 的组件作为节点的工作单元专注于解决如何将这些单元有机地、有状态地组织成一个稳健流程的问题。你可以只用 LangChain 来构建简单的链Chain但当你需要复杂、有状态、多步骤的交互时LangGraph 是更专业的选择。它们不是替代关系而是互补LangGraph 负责调度和状态LangChain 负责提供可调度的“零件”。2. 实战第一步设计状态与构建基础 StateGraph理论说再多不如动手。我们以一个“智能研究助手”为项目目标用户输入一个复杂问题智能体需要自动进行网络搜索、阅读相关资料、总结分析并最终生成一份结构化的报告。这个流程涉及多个步骤和决策点非常适合用 LangGraph 来构建。2.1 定义核心状态State状态是整个工作流的“中央数据总线”。我们需要仔细设计它包含哪些信息。一个好的状态设计应该包含输入、中间产物、最终输出以及控制标志。from typing import TypedDict, List, Annotated from langgraph.graph.message import add_messages import operator class GraphState(TypedDict): 定义工作流的状态结构。 # 用户输入 question: str # 控制流程的标志位例如当前阶段 current_step: str # 中间数据搜索查询词、搜索结果、提取的文本 search_queries: List[str] search_results: List[dict] # 可能包含标题、链接、摘要 extracted_texts: List[str] # 中间数据分析过程中的草稿、要点 analysis_draft: str key_points: List[str] # 最终输出 final_report: str # 消息历史如果涉及多轮对话或智能体间通信 messages: Annotated[list, add_messages]这里我们使用了TypedDict来获得类型提示。Annotated用于messages字段是 LangGraph 处理消息历史的一种便捷方式。current_step是一个简单的控制标志用于决定下一个节点。2.2 创建节点Node与边Edge节点是执行具体工作的函数。每个函数接收整个GraphState作为输入修改其中部分内容然后返回更新后的状态。我们先创建两个基础节点search_web和extract_content。from langchain_community.tools import TavilySearchResults from langchain_community.document_loaders import AsyncHtmlLoader from langchain_community.document_transformers import Html2TextTransformer # 初始化工具需要TAVILY_API_KEY search_tool TavilySearchResults(max_results3) def search_web_node(state: GraphState) - GraphState: 节点根据问题生成搜索词并执行搜索。 question state[question] # 1. 生成搜索查询词这里简化直接使用问题。进阶版可以用LLM优化查询词 queries [question] # 模拟可能生成多个查询词 # queries llm.invoke(f针对问题{question}生成3个不同的搜索关键词。) # 2. 执行搜索 all_results [] for query in queries: try: results search_tool.invoke(query) all_results.extend(results) except Exception as e: print(f搜索查询 {query} 时出错: {e}) # 可以在这里将错误信息存入state供后续节点处理 # 更新状态 return { search_queries: queries, search_results: all_results, current_step: content_extraction # 更新步骤标志 } def extract_content_node(state: GraphState) - GraphState: 节点从搜索结果中加载并提取纯文本内容。 results state.get(search_results, []) urls [res.get(url) for res in results if res.get(url)] if not urls: return {extracted_texts: [], current_step: analyze} # 使用异步加载器提高效率 loader AsyncHtmlLoader(urls) docs loader.load() # 将HTML转换为纯净文本 transformer Html2TextTransformer() transformed_docs transformer.transform_documents(docs) extracted_texts [doc.page_content[:5000] for doc in transformed_docs] # 截断避免过长 return { extracted_texts: extracted_texts, current_step: analyze }有了节点我们开始构建图。from langgraph.graph import StateGraph, END # 1. 创建图构建器 workflow StateGraph(GraphState) # 2. 添加节点 workflow.add_node(search_web, search_web_node) workflow.add_node(extract_content, extract_content_node) # 3. 添加边定义执行顺序 workflow.add_edge(search_web, extract_content) workflow.add_edge(extract_content, END) # 暂时直接结束 # 4. 设置入口点 workflow.set_entry_point(search_web) # 5. 编译图 app workflow.compile()现在一个最简单的两节点线性工作流就完成了。你可以运行它initial_state GraphState(question什么是LangGraph它和LangChain有什么区别, current_stepstart) final_state app.invoke(initial_state) print(final_state[extracted_texts][0][:500]) # 打印部分提取的文本这只是一个开始。目前的工作流是僵硬的直线。接下来我们要引入条件分支让它具备决策能力。3. 引入决策与循环让工作流“智能”起来一个只会机械执行固定步骤的流程不是智能体。我们需要它能够根据中间结果做出判断。比如如果搜索返回的结果太少或质量不高我们应该触发新一轮搜索如果提取的文本已经足够回答简单问题或许可以跳过深度分析。3.1 使用条件边Conditional Edge我们修改流程在extract_content节点后添加一个路由决策判断提取的文本是否足够进行下一步分析。首先创建一个路由函数。这个函数不修改状态只根据当前状态返回下一个节点的名称。def route_after_extraction(state: GraphState) - str: 路由函数判断内容提取后是进行分析还是重新搜索。 texts state.get(extracted_texts, []) # 简单的启发式规则如果没提取到文本或者文本总长度太短则重新搜索 if not texts or sum(len(t) for t in texts) 500: # 可以在这里修改搜索策略例如更新查询词 new_queries state.get(search_queries, []) [补充信息] state[search_queries] new_queries return search_web # 返回需要重新执行的节点名 else: return analyze # 进入分析节点然后修改我们的图构建逻辑用add_conditional_edges来替代之前直接指向END的边。# 首先我们需要一个分析节点暂未实现先作为占位符 def analyze_node(state: GraphState) - GraphState: 节点分析提取的文本生成要点。 # 这里后续会集成LLM state[key_points] [要点1, 要点2] state[current_step] generate_report return state def generate_report_node(state: GraphState) - GraphState: 节点根据要点生成最终报告。 # 这里后续会集成LLM state[final_report] 这是一份基于分析的模拟报告。 return state # 重建图 workflow StateGraph(GraphState) # 添加所有节点 workflow.add_node(search_web, search_web_node) workflow.add_node(extract_content, extract_content_node) workflow.add_node(analyze, analyze_node) workflow.add_node(generate_report, generate_report_node) # 设置线性边 workflow.add_edge(search_web, extract_content) workflow.add_edge(analyze, generate_report) workflow.add_edge(generate_report, END) # 关键为 extract_content 节点添加条件边 workflow.add_conditional_edges( extract_content, route_after_extraction, # 路由函数 { search_web: search_web, # 如果返回search_web则跳转到search_web节点 analyze: analyze # 如果返回analyze则跳转到analyze节点 } ) # 设置入口 workflow.set_entry_point(search_web) app workflow.compile()现在工作流具备了基本的“判断-循环”能力。如果第一次搜索和提取的内容不够它会尝试修改查询词这里简单追加了“补充信息”并重新搜索直到内容达标或达到隐式循环上限需要额外机制避免无限循环。这就是状态化工作流的威力状态search_queries在循环中被修改并影响下一轮执行。3.2 防止无限循环与超时处理在真实场景中必须避免无限循环。LangGraph 提供了interrupt和Timeout等机制但更常见的做法是在状态中设置一个计数器并在路由函数中检查。class GraphState(TypedDict): # ... 其他字段同上 ... search_attempt_count: int # 新增搜索尝试计数器 def route_after_extraction(state: GraphState) - str: texts state.get(extracted_texts, []) attempt state.get(search_attempt_count, 0) # 规则1如果尝试超过3次强制进入分析即使内容不足 if attempt 3: return analyze # 规则2内容不足且尝试次数未超限则重新搜索 if not texts or sum(len(t) for t in texts) 500: state[search_attempt_count] attempt 1 new_queries state.get(search_queries, []) [f补充信息尝试{attempt1}] state[search_queries] new_queries return search_web else: return analyze4. 从单智能体到多智能体协作当任务足够复杂时让一个“全能”的智能体做所有事情会导致其逻辑臃肿且容易出错。更好的模式是多智能体协作每个智能体对应一个或多个节点职责单一通过状态流进行通信和协作。在我们的研究助手项目中可以设计三个智能体研究员Researcher负责搜索和获取信息。对应search_web和extract_content节点。分析师Analyst负责阅读、理解和提炼信息要点。对应analyze节点。撰稿人Writer负责根据要点组织语言生成格式优美的最终报告。对应generate_report节点。4.1 为智能体注入“大脑”集成大模型之前的节点都是基于规则或简单工具。现在我们为分析师和撰稿人节点集成大模型如 OpenAI GPT-4、 Anthropic Claude 或本地模型让它们真正具备理解和生成能力。首先定义两个不同角色的 LLM 调用。在实践中可以为不同角色设定不同的系统提示词System Prompt甚至使用不同能力的模型。from langchain_openai import ChatOpenAI # 假设使用 OpenAI API llm_analyst ChatOpenAI(modelgpt-4-turbo-preview, temperature0.1) llm_writer ChatOpenAI(modelgpt-4-turbo-preview, temperature0.7) # 撰稿人可以更有创造性 def analyze_node(state: GraphState) - GraphState: 分析师智能体阅读文本提炼关键要点。 texts state.get(extracted_texts, []) question state[question] if not texts: state[key_points] [未找到相关信息。] return state # 构建给分析师的提示词 context \n\n---\n\n.join(texts[:3]) # 限制上下文长度 prompt f 你是一位严谨的研究分析师。请基于以下背景资料针对用户问题“{question}”提炼出3-5个最核心的要点。 要点要求客观、简洁、基于资料。 背景资料 {context} 核心要点每条用‘- ’开头 response llm_analyst.invoke(prompt) points response.content.strip().split(\n) # 简单清理 points [p.strip(- ).strip() for p in points if p.strip().startswith(-)] state[key_points] points if points else [分析未能提取出明确要点。] state[current_step] generate_report return state def generate_report_node(state: GraphState) - GraphState: 撰稿人智能体根据要点撰写结构化报告。 question state[question] key_points state.get(key_points, []) prompt f 你是一位专业的科技报告撰稿人。请根据以下要点围绕问题“{question}”撰写一份结构清晰、语言流畅的简短报告约300字。 报告结构建议 1. 引言重申问题。 2. 核心发现基于要点展开。 3. 总结与意义。 要点 {chr(10).join(f- {p} for p in key_points)} 报告 response llm_writer.invoke(prompt) state[final_report] response.content return state现在我们的工作流拥有了两个具备 LLM 能力的智能体节点。它们通过共享的GraphState特别是extracted_texts和key_points字段进行协作。研究员为分析师准备材料分析师为撰稿人提炼骨架撰稿人最终成文。这种基于状态共享的协作模式比让单个 LLM 一次性完成所有步骤更加可控、可解释也更容易针对单个环节进行优化或替换模型。5. 集成 RAG为智能体注入“长期记忆”与“领域知识”多智能体协作解决了流程问题但它们的知识仍然局限于单次会话和通用模型。对于专业领域问题如公司内部文档、技术手册、行业报告我们需要为智能体配备RAG检索增强生成能力使其能访问外部知识库。RAG 不是替代工作流而是增强其中某个或某几个智能体的能力。例如可以让“研究员”智能体不仅搜索公开网络还能检索内部知识库或者让“分析师”智能体在分析时结合检索到的内部资料进行交叉验证。5.1 在状态流中嵌入 RAG 检索节点假设我们已经有一个构建好的向量数据库例如使用 Chroma、Weaviate 或 Pinecone里面存储了公司内部技术文档。我们创建一个新的节点retrieve_internal_knowledge并将其插入到工作流中合适的位置。from langchain_chroma import Chroma from langchain_openai import OpenAIEmbeddings from langchain.text_splitter import RecursiveCharacterTextSplitter # 注意以下为示例你需要预先准备好向量库 # 初始化嵌入模型和向量库连接 embeddings OpenAIEmbeddings(modeltext-embedding-3-small) vectorstore Chroma(persist_directory./my_chroma_db, embedding_functionembeddings) retriever vectorstore.as_retriever(search_kwargs{k: 3}) # 检索最相关的3条 def retrieve_knowledge_node(state: GraphState) - GraphState: 节点从内部知识库检索相关信息。 question state[question] # 执行检索 relevant_docs retriever.invoke(question) # 将检索到的文档内容存入状态供后续节点使用 internal_knowledge [doc.page_content for doc in relevant_docs] # 可以单独存放也可以合并到之前的 extracted_texts 中 state[internal_knowledge] internal_knowledge state[current_step] analyze # 假设检索后直接进入分析 return state然后修改工作流图在搜索公开信息后并行或串行地检索内部知识。# 修改图在 extract_content 后analyze 前加入知识检索 workflow StateGraph(GraphState) workflow.add_node(search_web, search_web_node) workflow.add_node(extract_content, extract_content_node) workflow.add_node(retrieve_knowledge, retrieve_knowledge_node) workflow.add_node(analyze, analyze_node) workflow.add_node(generate_report, generate_report_node) # 线性流程搜索 - 提取 - 检索 - 分析 - 报告 workflow.add_edge(search_web, extract_content) workflow.add_edge(extract_content, retrieve_knowledge) workflow.add_edge(retrieve_knowledge, analyze) workflow.add_edge(analyze, generate_report) workflow.add_edge(generate_report, END) workflow.set_entry_point(search_web) app workflow.compile()现在分析师智能体在生成要点时就可以同时参考extracted_texts来自网络和internal_knowledge来自内部知识库。我们需要修改analyze_node的提示词将两部分信息都整合进去。def analyze_node(state: GraphState) - GraphState: texts state.get(extracted_texts, []) internal_knowledge state.get(internal_knowledge, []) question state[question] # 合并上下文 all_context_parts [] if texts: all_context_parts.append(【来自公开网络的信息】\n \n\n---\n\n.join(texts[:3])) if internal_knowledge: all_context_parts.append(【来自内部知识库的信息】\n \n\n---\n\n.join(internal_knowledge[:3])) if not all_context_parts: state[key_points] [未找到任何相关信息。] return state context \n\n.join(all_context_parts) prompt f 你是一位严谨的研究分析师。请综合以下来自公开网络和内部知识库的资料针对用户问题“{question}”提炼出3-5个最核心的要点。 注意比较和交叉验证不同来源的信息。 背景资料 {context} 核心要点每条用‘- ’开头 # ... 后续调用LLM的代码不变 ...通过这种方式RAG 被无缝地编织进了智能体的决策流程中成为其“长期记忆”和“领域知识”的来源。这比单纯用一个 RAG 问答链要强大得多因为检索到的知识是在一个更复杂、多步骤的推理流程中被使用的。5.2 多模型接入为不同任务选择最合适的“大脑”在上面的例子中分析师和撰稿人都用了 GPT-4。但在实际生产中成本、速度和任务特性都需要考虑。LangChain/LangGraph 的另一个优势是模型无关性。你可以轻松地为不同节点切换不同的模型提供商。例如可以让“研究员”使用快速且便宜的模型如 GPT-3.5-Turbo来生成搜索查询词让“分析师”使用能力强、上下文窗口大的模型如 Claude-3-Opus 或 GPT-4进行深度分析让“撰稿人”使用擅长创意写作的模型比如特定调优过的模型。只需在初始化节点时传入不同的 LLM 实例即可。LangChain 的统一接口让这变得非常简单。from langchain_anthropic import ChatAnthropic from langchain_google_genai import ChatGoogleGenerativeAI # 为不同智能体配置不同模型 llm_researcher ChatOpenAI(modelgpt-3.5-turbo, temperature0) # 便宜用于生成查询词 llm_analyst ChatAnthropic(modelclaude-3-opus-20240229, temperature0.1) # 能力强用于分析 llm_writer ChatGoogleGenerativeAI(modelgemini-pro, temperature0.7) # 流畅用于写作 # 然后在对应的节点函数中使用各自的 llm 实例这种灵活性使得你可以构建一个“混合模型”的智能体系统在效果、成本和速度之间取得最佳平衡。6. 工业级落地的关键可观测性、持久化与错误处理一个能在生产环境跑起来的智能体光有核心逻辑是不够的。它必须健壮、可监控、可调试。6.1 可观测性记录每一步的输入输出LangGraph 提供了Checkpointer和Message等机制来持久化状态但对于日志我们通常需要更灵活的方式。一个简单有效的方法是在每个节点的函数内部将关键操作和结果记录到结构化日志中。import logging import json logger logging.getLogger(__name__) def search_web_node(state: GraphState) - GraphState: question state[question] # 记录输入 logger.info(f[search_web_node] 开始处理问题: {question}, extra{state: state}) # ... 执行搜索 ... # 记录输出和关键结果 logger.info(f[search_web_node] 生成查询词: {queries}, extra{queries: queries}) logger.info(f[search_web_node] 获得结果数: {len(all_results)}, extra{sample_result: all_results[0] if all_results else None}) # 对于错误记录异常 except Exception as e: logger.error(f[search_web_node] 搜索失败: {e}, exc_infoTrue) # 可以选择将错误信息存入state让后续节点或路由函数处理 state[last_error] str(e) return new_state将这些日志收集到如 ELK、Loki 或 Datadog 等系统中你就能清晰地看到每个请求在智能体工作流中的完整生命周期便于排查问题和分析性能瓶颈。6.2 持久化与断点续跑使用 Checkpoints对于长时间运行或需要中断后恢复的工作流LangGraph 的 Checkpoint 机制至关重要。它允许你在每个节点执行后将完整的GraphState持久化到数据库如 Redis、PostgreSQL。当系统重启或从故障中恢复时可以从上一个检查点继续执行。from langgraph.checkpoint.sqlite import SqliteSaver from langgraph.graph import StateGraph # 初始化一个 SQLite 检查点存储器 memory SqliteSaver.from_conn_string(:memory:) # 生产环境用文件或网络数据库 workflow StateGraph(GraphState, config_schemaMyConfigSchema) # 需要定义ConfigSchema # ... 添加节点和边 ... app workflow.compile(checkpointermemory) # 编译时传入 checkpointer # 调用时需要传入一个线程IDthread_id用于标识同一个会话或流程 config {configurable: {thread_id: user_123_session_456}} initial_state GraphState(...) # invoke 会返回最终状态并在后台自动保存检查点 final_state app.invoke(initial_state, configconfig) # 如果中途中断可以通过 get_state 获取最新状态并继续 invoke state app.get_state(config) # app.invoke(None, config) 可以从该状态继续执行如果图支持从中间开始6.3 系统化的错误处理错误处理不能只靠try...except。需要在状态图和节点设计层面考虑节点级重试对于网络请求如搜索、LLM调用等暂时性错误可以在节点内部实现重试逻辑。工作流级容错通过条件边将错误信息state[“last_error”]作为路由依据。例如如果搜索节点连续失败可以路由到一个“人工降级”节点或直接生成一个包含错误说明的最终报告。超时控制为耗时长的节点如 LLM 调用设置超时防止整个工作流卡死。这通常需要结合异步和外部超时机制。降级方案当核心组件如某个 LLM API不可用时是否有备选模型或规则兜底7. 总结从项目到平台的思维转变通过这个完整的实战拆解我们可以看到使用 LangGraph LangChain 构建工业级 Agent远不止是调用 API 的简单叠加。它要求我们完成一次思维转变从编写线性脚本到设计状态化工作流思考的不再是一行行代码而是由节点、边和状态构成的可视化流程图。从打造“全能模型”到构建“协作团队”将复杂任务拆解由多个职责单一的智能体节点通过状态共享协同完成系统更健壮也更易优化。从单一模型调用到混合模型编排根据子任务的特点灵活选用不同能力、成本、速度的模型实现性价比最大化。从演示原型到可观测系统将日志、检查点、错误处理、降级方案作为一等公民来设计确保系统在生产环境中可靠运行。最终你构建的不是一个“智能体”而是一个智能体框架或平台。新的任务可以通过设计新的状态图和组装不同的节点复用或新建来快速实现。这才是 LangGraph 带来的长期价值它提供了一套用于构建复杂、可靠、可维护的 AI 应用的标准范式和解耦架构。当你下次再面对一个需要多步骤、多判断、多数据源的自动化任务时不妨先拿起笔画一画它的状态转移图。你会发现很多复杂性问题在清晰的架构面前都会迎刃而解。