ARTICLE DETAIL

资讯详情

深耕网站SEO优化与搜索引擎排名提升的一线实战洞察。

LangGraph进阶:状态持久化、人工介入与多智能体协作实战

LangGraph进阶:状态持久化、人工介入与多智能体协作实战 1. 从“线性脚本”到“状态驱动”为什么需要进阶的工作流如果你用过 LangChain 或者自己写过一些简单的 AI 应用脚本大概会经历这样一个过程写一个函数调用大模型 API处理返回结果再根据结果决定下一步。这就像写一个线性的剧本角色Agent按部就班地说台词。但当剧情变得复杂角色增多甚至需要根据观众用户的实时反馈来调整剧情时这种“剧本式”的编程就捉襟见肘了。这就是LangGraph要解决的核心问题。它不是一个替代 LangChain 的工具而是一个基于 LangChain 构建的、专门用于编排复杂、有状态、多步骤 AI 工作流的框架。你可以把它想象成一个可视化、可编程的流程图引擎专门为 AI 智能体Agent设计。我们之前可能已经了解了如何用 LangGraph 构建一个简单的对话链或工具调用流程但那只是入门。当我们要构建真正可用于生产环境的系统时三个进阶能力至关重要状态持久化检查点、人工介入Human-in-the-loop和多智能体协作。想象一下这些场景一个自动化的客户支持系统在处理到一半时需要转接给人工客服确认敏感信息一个多步骤的数据分析流水线其中某个步骤失败后我们希望能从失败点继续而不是重头开始一个由“研究员”、“写手”、“审阅员”多个 AI 角色组成的团队如何高效协作完成一份报告这些需求正是 LangGraph 进阶功能大显身手的地方。本文将深入这三个核心进阶主题结合代码示例和设计思路帮你把 LangGraph 从“玩具”升级为“生产级工具”。2. 状态持久化理解与应用检查点Checkpoint在简单的工作流中状态State是内存中的临时对象流程跑完就消失了。但对于一个可能运行数小时、包含昂贵模型调用比如 GPT-4的复杂工作流来说这是不可接受的。我们需要检查点Checkpoint机制将工作流在任意节点的完整状态保存下来允许系统在中断如程序崩溃、服务器重启后从中断点恢复也支持暂停、继续等操作。2.1 检查点的核心概念与价值LangGraph 中的检查点不仅仅是保存一个变量值。它保存的是整个State对象的序列化快照包括当前节点Node工作流执行到了哪个步骤。状态值ValuesState对象中所有键的值。下一步Next根据当前状态和节点逻辑计算出的下一个要执行的节点或节点列表。元数据Metadata可选可以包含时间戳、版本、创建者等信息。其核心价值在于容错与恢复进程崩溃或部署更新后可以从最新的检查点继续避免重复执行已完成的昂贵步骤如大模型调用。调试与审计可以回溯工作流的历史状态查看每一步的输入输出便于定位问题。支持异步与长时运行工作流可以暂停等待外部事件如人工审核、第三方API回调事件到来后从检查点继续。实现“续跑”功能用户关闭了网页下次打开可以继续之前的会话。2.2 如何实现检查点配置持久化存储LangGraph 通过CheckpointSaver这一抽象来实现持久化。它不是一个开箱即用的数据库而是一个接口你需要为其提供一个具体的存储后端。社区中已有一些实现你也可以自己实现。一个常见的搭配是使用SqliteSaver它轻量且易于集成。下面是一个完整的示例展示如何创建一个带检查点的工作流。首先定义我们的状态。假设我们有一个文档总结和提问的工作流。from typing import TypedDict, List, Annotated from langgraph.graph import StateGraph, END from langgraph.checkpoint.sqlite import SqliteSaver import operator # 1. 定义状态结构 class State(TypedDict): # 原始文档 original_document: str # 生成的摘要 summary: str # 根据摘要生成的问题列表 questions: List[str] # 用户对问题的答案 answers: Annotated[List[str], operator.add] # 使用operator.add来追加列表 # 当前步骤 step: str # 2. 创建持久化存储后端 # 这里使用SQLite数据会保存在 checkpoints.db 文件中 memory SqliteSaver.from_conn_string(checkpoints.db) # 3. 构建图并传入 checkpointer builder StateGraph(State, config_schemadict) # config_schema 用于区分不同工作流实例 builder.add_node(summarize, summarize_node) builder.add_node(generate_questions, generate_questions_node) builder.add_node(collect_answers, collect_answers_node) builder.set_entry_point(summarize) builder.add_edge(summarize, generate_questions) builder.add_edge(generate_questions, collect_answers) builder.add_edge(collect_answers, END) # 关键步骤将 memory 作为 checkpointer 传入 graph builder.compile(checkpointermemory)现在当我们运行这个图时需要提供一个config参数其中包含configurable字段来指定本次运行的thread_id。thread_id是唯一标识一个工作流会话的键所有检查点都会通过它来存储和读取。# 运行工作流并指定 thread_id config {configurable: {thread_id: user_session_12345}} initial_state {original_document: 这里是一篇很长的技术文章...} # 第一次运行 result graph.invoke(initial_state, configconfig) print(result[summary]) # 输出摘要 # 模拟中断程序在这里结束 # 第二次运行恢复我们只需要相同的 config并从空状态或部分状态开始调用。 # LangGraph 会自动加载最新的检查点并从上次中断的节点继续执行。 new_result graph.invoke({}, configconfig) # 注意这里传入空状态或仅更新部分状态 print(new_result[questions]) # 会输出之前 generate_questions 节点产生的问题注意在恢复执行时invoke传入的初始状态会与从检查点加载的状态进行合并。通常如果你只是想继续传入空字典{}即可。如果你想用新数据覆盖状态的某些部分例如更新了文档你可以传入对应的键值对。2.3 检查点的高级用法与避坑指南1. 并发安全与锁机制 当多个进程或线程同时尝试恢复和更新同一个thread_id的检查点时会发生竞争条件。SqliteSaver利用数据库事务提供了一定的并发安全。但对于高并发场景你可能需要在应用层或使用支持乐观锁/悲观锁的存储后端如 Redis、PostgreSQL来实现更精细的控制。2. 状态序列化与版本控制 如果你的State结构即TypedDict的字段发生了变更比如新增或删除了一个字段旧的检查点在反序列化时可能会失败。在生产环境中你需要考虑状态模式的版本迁移策略。一个简单的方法是在State中保留一个version字段并在加载旧检查点时编写迁移逻辑。3. 存储成本与清理 检查点会不断创建。对于一个长期运行、频繁触发的工作流存储空间可能快速增长。你需要定期清理旧的检查点。SqliteSaver本身不提供自动清理你需要自己执行 SQL 语句如DELETE FROM checkpoints WHERE thread_id ? AND timestamp ?或使用其他存储后端的 TTL生存时间功能。4. 配置化Configurable的威力config参数中的configurable字段非常强大。除了thread_id你还可以存储其他会话级元数据比如user_id、project_id等。这些信息也会被保存在检查点中并在恢复时可用。这允许你基于更丰富的上下文来恢复工作流。config { configurable: { thread_id: analysis_789, user_id: user_abc, priority: high } }3. 人工介入实现 Human-in-the-Loop 工作流完全自动化的 AI 工作流虽然高效但在处理关键决策、创造性任务或涉及伦理、安全的场景时引入人工判断是必要且明智的。LangGraph 通过“暂停”并等待外部输入的机制来优雅地支持这一点。3.1 核心机制interrupt_before与interrupt_afterLangGraph 允许你在指定的节点之前或之后设置中断点。当工作流执行到这些点时它会自动暂停将控制权交还给调用者并返回一个特殊的结果表明它正在等待某个特定节点的输入。interrupt_before[“node_name”]: 在进入node_name节点之前暂停。这通常用于在AI执行前由人来提供输入或指令。interrupt_after[“node_name”]: 在离开node_name节点之后暂停。这通常用于在AI产生输出后由人来审核、修改或确认结果。3.2 实战构建一个带人工审核的文档发布流程让我们构建一个简单的博客发布工作流AI 生成初稿 - 人工审核 - AI 根据反馈修改 - 最终发布。from typing import Literal from langgraph.graph import StateGraph, END, MessagesState from langgraph.prebuilt import ToolNode from langgraph.checkpoint.sqlite import SqliteSaver # 使用 MessagesState 简化对话状态管理 class State(MessagesState): draft: str human_feedback: str final_content: str status: Literal[drafting, awaiting_review, revising, published] drafting # 定义节点函数 def write_draft_node(state: State): # 模拟AI写作 new_draft f基于最新趋势这是一篇关于LangGraph的博文草稿。当前状态: {state[status]} return {draft: new_draft, status: awaiting_review} def revise_draft_node(state: State): # 模拟AI根据反馈修改 feedback state.get(human_feedback, 无反馈) revised f{state[draft]}\n\n---\n已根据反馈 {feedback} 进行修改。 return {draft: revised, status: revising} def publish_node(state: State): return {final_content: state[draft], status: published} # 构建图 builder StateGraph(State) builder.add_node(write_draft, write_draft_node) builder.add_node(revise_draft, revise_draft_node) builder.add_node(publish, publish_node) builder.set_entry_point(write_draft) # 关键在 write_draft 之后设置中断等待人工审核 builder.add_edge(write_draft, revise_draft) builder.add_edge(revise_draft, publish) builder.add_edge(publish, END) memory SqliteSaver.from_conn_string(:memory:) # 配置中断在 ‘revise_draft’ 节点之前中断因为我们需要在AI修改前拿到人工反馈。 graph builder.compile( checkpointermemory, interrupt_before[revise_draft] # 这里设置中断点 )运行这个工作流config {configurable: {thread_id: blog_post_1}} # 1. 首次调用AI写草稿然后在 revise_draft 前中断 result graph.invoke({status: drafting}, configconfig) print(result) # 输出可能包含一个特殊标记或者你可以检查状态。实际上invoke 会抛出一个 Interruption 异常或返回特定结构。 # 在 LangGraph 的流式或事件处理中更常见的模式是使用 stream 并监听事件。 # 为了清晰我们换一种方式演示使用 stream 并处理 Send 事件。 from langgraph.graph import Send # 假设我们有一个处理函数 def run_workflow_with_human(): thread_config {configurable: {thread_id: blog_post_2}} # 初始化 for event in graph.stream({status: drafting}, thread_config, stream_modevalues): if isinstance(event, dict): print(f状态更新: {event}) # 在实际应用中这里会有一个事件循环当检测到需要人工介入时就跳出循环或发送通知。 # 更直观的 API 使用get_state 和 update_state # 首先运行到中断点 try: graph.invoke({status: drafting}, configconfig) except Exception as e: # 这里会捕获到中断在实际框架中可能有更优雅的交互方式。 # 例如Dify、Coze 等平台封装了此过程在UI上生成一个“待办事项”。 pass # 此时工作流已暂停。我们可以获取当前状态。 current_state graph.get_state(config) print(f当前草稿: {current_state.values[draft]}) print(f当前状态: {current_state.values[status]}) # 应该是 awaiting_review # 现在模拟人工审核提供反馈 human_feedback_input 开头不够吸引人请加入一个具体的案例。 # 将反馈更新到状态中并告诉图继续执行从中断点 revise_draft 开始 updated_state {human_feedback: human_feedback_input} graph.update_state(config, updated_state) # 继续执行从 revise_draft 节点开始 result graph.invoke({}, configconfig) # 传入空状态或仅包含更新的状态 print(f最终内容: {result[final_content]}) print(f发布状态: {result[status]})实操心得在真实的后端服务中你通常不会用try...except来捕获中断。更好的模式是使用graph.stream()异步执行工作流。监听事件类型当遇到Send事件其name字段指向一个中断的节点时将工作流state和config存入数据库并生成一个任务ID。暴露一个 REST API如POST /task/{task_id}/feedback供前端或人工审核界面调用。当API接收到反馈后从数据库加载状态和配置调用graph.update_state()和graph.invoke()继续执行。这种模式天然与检查点机制结合非常健壮。3.3 设计人工介入点的考量中断的粒度是在一个复杂节点的前后中断还是将复杂节点拆分成多个小节点在中间中断后者更灵活但图结构会更复杂。超时与备选方案如果人工审核迟迟没有响应怎么办你需要设置超时机制并在工作流中设计“默认路径”或“升级路径”例如超时后自动发送提醒或转给另一位审核者。反馈的结构化尽量让人工提供结构化的反馈如从下拉框选择“通过”、“拒绝并重写”、“需修改XX部分”而不是纯文本。这能简化后续AI处理反馈的逻辑。你可以在State中定义专门的字段来存储结构化反馈。4. 多智能体协作构建角色化团队工作流单个“全能”的 Agent 往往力有不逮。更强大的模式是模拟一个团队让多个各司其职的 Agent 协作完成任务。LangGraph 的图结构非常适合对这种协作关系进行建模。4.1 设计模式管理者与工作者一种常见的设计模式是“管理者Supervisor-工作者Worker”。管理者 Agent负责理解总体任务进行任务规划与分解并将子任务分配给特定的工作者最后汇总和评估结果。工作者 Agent是领域专家负责执行具体的子任务如编写代码、检索信息、分析数据等。在 LangGraph 中管理者和工作者都可以实现为一个节点Node。管理者节点根据当前状态和任务决定下一步调用哪个工作者节点这通过条件边Conditional Edge来实现。4.2 实战构建一个技术调研团队假设我们要组建一个团队来调研“向量数据库的最新进展”团队包括规划师Planner分解调研任务。研究员Researcher使用网络搜索工具查找信息。分析师Analyst对搜集的信息进行归纳总结。写手Writer将分析结果整理成格式良好的报告。from typing import List from langchain_core.messages import HumanMessage, SystemMessage from langchain_openai import ChatOpenAI from langchain_community.tools import DuckDuckGoSearchRun from langgraph.graph import StateGraph, END # 定义状态 class ResearchState(TypedDict): original_query: str # 原始问题 plan: List[str] # 调研计划 research_results: List[str] # 搜集到的资料 analysis: str # 分析结论 report: str # 最终报告 current_role: str # 当前正在执行的角色 max_turns: int 10 # 最大循环次数防止无限循环 # 初始化工具和模型 search DuckDuckGoSearchRun() llm ChatOpenAI(modelgpt-4-turbo) # 1. 规划师节点 def planner_node(state: ResearchState): query state[original_query] prompt f 你是一个资深技术调研规划师。请将以下调研任务分解为3-5个具体的、可执行的子任务。 任务{query} 请以清晰的列表形式返回子任务。 messages [SystemMessage(content你是一个高效的任务规划师。), HumanMessage(contentprompt)] response llm.invoke(messages) # 假设返回格式是 “1. ...\n2. ...” plan [line.strip() for line in response.content.split(\n) if line.strip().startswith((1., 2., 3., 4., 5., -))] return {plan: plan, current_role: planner} # 2. 研究员节点 def researcher_node(state: ResearchState): # 从计划中取出第一个未完成的调研点这里简化处理每次执行一个 # 实际中状态里可能需要一个 current_plan_index if not state.get(plan): return {research_results: [], current_role: researcher} # 简化用第一个计划项去搜索 search_query state[plan][0] 最新进展 2024 result search.run(search_query) # 将结果累积到 research_results 中 current_results state.get(research_results, []) current_results.append(f调研点{state[plan][0]}\n结果{result[:500]}...) # 截断 return {research_results: current_results, current_role: researcher} # 3. 分析师节点 def analyst_node(state: ResearchState): all_results \n---\n.join(state.get(research_results, [])) prompt f 你是一个技术分析师。请根据以下搜集到的资料总结关于“{state[original_query]}”的核心观点、技术对比和趋势。 资料 {all_results} 请给出结构化的分析摘要。 messages [SystemMessage(content你是一个洞察力强的技术分析师。), HumanMessage(contentprompt)] response llm.invoke(messages) return {analysis: response.content, current_role: analyst} # 4. 写手节点 def writer_node(state: ResearchState): analysis state.get(analysis, ) prompt f 你是一名技术文档写手。请将以下分析内容整理成一篇适合发布在技术博客上的简短报告。 要求有标题、引言、核心发现分点论述、总结。 分析内容 {analysis} messages [SystemMessage(content你是一名优秀的科技文章写手。), HumanMessage(contentprompt)] response llm.invoke(messages) return {report: response.content, current_role: writer} # 5. 路由逻辑管理者逻辑 # 这是一个决定下一个节点是谁的函数 def route_after_planner(state: ResearchState) - str: # 规划完成后如果有计划项就去研究否则直接分析可能不需要研究 if state.get(plan): return researcher else: return analyst def route_after_researcher(state: ResearchState) - str: # 研究员完成后检查是否还有未完成的计划项这里简化只执行一次研究 # 实际中这里可以更复杂比如循环执行直到计划完成。 # 现在我们假设研究一次后就去分析。 return analyst def route_after_analyst(state: ResearchState) - str: # 分析完成后交给写手 return writer def route_after_writer(state: ResearchState) - str: # 写手完成后结束 return END # 构建图 builder StateGraph(ResearchState) builder.add_node(planner, planner_node) builder.add_node(researcher, researcher_node) builder.add_node(analyst, analyst_node) builder.add_node(writer, writer_node) builder.set_entry_point(planner) # 使用条件边进行路由 builder.add_conditional_edges( planner, route_after_planner, {researcher: researcher, analyst: analyst} # 路由函数返回字符串映射到节点名 ) builder.add_conditional_edges( researcher, route_after_researcher, {analyst: analyst} ) builder.add_conditional_edges( analyst, route_after_analyst, {writer: writer} ) builder.add_edge(writer, END) # 写手之后固定结束 graph builder.compile() # 运行工作流 initial_state {original_query: 向量数据库的最新技术进展, max_turns: 10} final_state graph.invoke(initial_state) print( 最终报告 ) print(final_state[report])4.3 多 Agent 协作的挑战与优化1. 共享上下文与信息隔离 所有 Agent 共享同一个State。这有利于信息传递但也可能导致“信息过载”或“意外修改”。好的实践是在State中为每个 Agent 设计清晰的命名空间例如researcher_findings,analyst_insights。使用Annotated类型和operator.add来安全地追加列表而不是覆盖。对于敏感信息可以考虑在子图中进行隔离。2. 循环与终止条件 多 Agent 协作容易陷入循环例如研究员和分析师互相要求对方提供更多信息。必须设置明确的终止条件在State中设置max_turns或max_iterations计数器。在路由逻辑中检查任务是否已完成例如所有计划项是否都已处理。使用“投票”或“共识”机制当多个 Agent 认为可以结束时才结束。3. 子图Subgraph封装复杂逻辑 如果一个工作者 Agent 的内部逻辑非常复杂例如它自己又是一个包含多步骤的规划-执行循环你可以将其封装成一个子图。主图只调用这个子图节点子图内部可以有自己的状态和逻辑。这极大地提升了模块化和可维护性。使用add_node时传入一个已编译的Graph对象即可。# 假设我们有一个复杂的代码生成子图 code_gen_graph ... # 另一个编译好的StateGraph builder StateGraph(MainState) # 将子图作为一个节点加入 builder.add_node(code_generator, code_gen_graph)4. 异步与并行执行 上面的例子是顺序执行的。LangGraph 支持基于状态键更新的并行执行。如果researcher和analyst节点不依赖于对方的输出即它们修改的是State中不同的键你可以通过配置让它们同时运行这在处理独立子任务时能大幅提升效率。这需要在定义节点和边时进行精细设计。5. 进阶整合构建一个带检查点、人工介入与多 Agent 的生产级工作流让我们将前面三个概念整合起来设计一个更贴近实际需求的场景一个自动化代码审查与合并工作流。需求描述当有新的 Pull Request (PR) 时自动触发工作流。代码分析 Agent检查代码风格和潜在 Bug。测试生成 Agent尝试为变更生成单元测试。安全扫描 Agent进行基础的安全漏洞检查。将三个 Agent 的结果汇总生成一份审查报告。如果报告中发现关键问题如高严重性安全漏洞则暂停工作流等待项目负责人人工确认。负责人可以给出“驳回”、“忽略并继续”、“添加评论”等指令。根据人工指令或自动判断工作流决定是自动合并、关闭 PR还是添加评论后等待。整个流程需要支持中断恢复比如负责人可能几天后才处理。这个工作流涵盖了多 Agent 协作分析、测试、安全、人工介入负责人审核和状态持久化支持长时间暂停。其 LangGraph 实现的核心骨架如下# 状态设计 class CodeReviewState(TypedDict): pr_id: str code_diff: str style_issues: List[str] generated_tests: str security_findings: List[dict] # 每个finding包含 {“level”: “high/medium/low“, “description”: “...”} review_report: str human_decision: Literal[approve, reject, comment, None] None human_comment: str final_action: Literal[merge, close, comment, None] None status: str “initialized” # 图结构概览 (伪代码) builder StateGraph(CodeReviewState) builder.add_node(“analyze_code”, analyze_code_node) # 代码分析Agent builder.add_node(“generate_tests”, generate_tests_node) # 测试生成Agent builder.add_node(“scan_security”, scan_security_node) # 安全扫描Agent builder.add_node(“compile_report”, compile_report_node) # 报告汇总Agent builder.add_node(“wait_for_human”, wait_for_human_node) # 人工介入节点可能只是一个标记 builder.add_node(“execute_decision”, execute_decision_node) # 执行最终动作的Agent builder.set_entry_point(“analyze_code”) # 并行执行分析、测试、安全扫描如果它们更新不同的状态键 builder.add_edge(“analyze_code”, “generate_tests”) builder.add_edge(“generate_tests”, “scan_security”) # 或者设计成真正的并行需要更精细的状态键配置 builder.add_edge(“scan_security”, “compile_report”) # 条件边根据报告决定是否需要人工介入 def after_report(state): has_critical_issue any(f[“level”] “high” for f in state.get(“security_findings”, [])) if has_critical_issue: return “wait_for_human” # 跳转到人工介入节点 else: return “execute_decision” # 自动执行 builder.add_conditional_edges( “compile_report”, after_report, {“wait_for_human”: “wait_for_human”, “execute_decision”: “execute_decision”} ) # 人工介入节点这里不执行实际逻辑只是设置一个中断点。 # 实际的人工交互在外部系统如Webhook处理。 def wait_for_human_node(state): # 这个节点可以什么都不做或者只是更新状态表示正在等待 return {“status”: “awaiting_human_review”} # 在 ‘wait_for_human’ 节点之后设置中断 # 当执行到这里图会暂停等待外部系统调用 update_state 来提供 human_decision 和 human_comment builder.add_edge(“wait_for_human”, “execute_decision”) # 最终决策节点 def execute_decision_node(state): decision state.get(“human_decision”) if decision “approve” or (decision is None and not state.get(“has_critical_issue”)): # 调用GitHub API合并PR action “merge” elif decision “reject”: # 调用GitHub API关闭PR action “close” else: # comment 或其他情况 # 调用GitHub API添加评论 action “comment” return {“final_action”: action, “status”: “completed”} builder.add_edge(“execute_decision”, END) # 编译图并启用 SQLite 检查点 memory SqliteSaver.from_conn_string(“code_review.db”) graph builder.compile( checkpointermemory, interrupt_after[“wait_for_human”] # 在等待人工节点后中断 )这个工作流可以通过一个 Web 服务器来驱动GitHub Webhook 收到 PR 事件触发graph.invoke(initial_state, config{“thread_id”: pr_id})。工作流运行到wait_for_human节点后中断状态被保存。服务器向项目负责人的聊天工具如 Slack发送通知并提供一个审批链接。负责人点击链接在网页上看到审查报告并做出决定。网页后端调用graph.update_state(config, {“human_decision”: “approve”, “human_comment”: “LGTM”})然后再次graph.invoke({}, config)。工作流从中断处继续执行execute_decision_node完成合并操作。通过这种方式我们利用 LangGraph 构建了一个鲁棒、可持久化、支持人机交互、由多个智能体协同的自动化流程这正是现代 AI 应用工程化的核心所在。
返回列表