09-Agent应用与工作流编排框架LangGraph
💡 学习目标
- 掌握 LangGraph 的核心概念和使用场景
- 掌握 LangGraph 构建 ReAct Agent
- 掌握 LangGraph 创建自定义工作流
- 掌握 Memory 与 Persistence
1. LangGraph 介绍
1.1 基本概述
LangGraph 是由 LangChain 团队开发的一个开源框架,旨在帮助开发者构建基于大型语言模型(LLM)的复杂、有状态、多主体的应用。它通过将工作流表示为图结构(graph),提供了更高的灵活性和控制能力,特别适合需要循环逻辑、状态管理以及多主体协作的场景,比如智能代理(agent)和多代理工作流。
LangGraph 是为智能体和工作流设计的一套底层编排框架,旨在构建、部署和管理复杂的生成式 AI 代理工作流。它提供了一套工具和库,使用户能够以可扩展且高效的方式创建、运行和优化大型语言模型(LLM)。LangGraph 的核心是利用基于图的架构的强大功能来建模和管理AI 代理工作流中各个组件之间的复杂关系。
官方文档:https://langchain-ai.github.io/langgraph/
1.2 核心概念
图结构(Graph Structure)
LangGraph 将应用逻辑组织成一个有向图,其中:
- 节点(Nodes):代表具体的操作或计算步骤,可以是调用语言模型、执行函数或与外部工具交互等
- 边(Edges):定义节点之间的连接和执行顺序,支持普通边(直接连接)和条件边(基于条件动态选择下一步)
状态管理(State Management)
LangGraph 的核心特点是自动维护和管理状态
状态(State)是一个贯穿整个图的共享数据结构,记录了应用运行过程中的上下文信息
每个节点可以根据当前状态执行任务并更新状态,确保系统在多步骤或多主体交互中保持一致性
循环能力(Cyclical Workflows)
与传统的线性工作流(如 LangChain 的 LCEL)不同,LangGraph 支持循环逻辑,这使得它非常适合需要反复推理、决策或与用户交互的代理应用。例如,一个代理可以在循环中不断调用语言模型,直到达成目标。
1.3 主要特点
灵活性: 开发者可以精细控制工作流的逻辑和状态更新,适应复杂的业务需求
持久性: 内置支持状态的保存和恢复,便于错误恢复和长时间运行的任务
多主体协作: 允许多个代理协同工作,每个代理负责特定任务,通过图结构协调交互
工具集成: 可以轻松集成外部工具(如搜索API)或自定义函数,增强代理能力
人性化交互: 支持“人机交互”(human-in-the-loop)功能,让人类在关键步骤参与决策
1.4 使用场景
LangGraph 特别适用于以下场景:
对话代理: 构建能够记住上下文、动态调整策略的智能聊天机器人
多步骤任务: 处理需要分解为多个阶段的复杂问题,如研究、写作或数据分析
多代理系统: 协调多个代理分工合作,比如一个负责搜索信息、另一个负责总结内容的系统
1.5 与 LangChain 的关系
- LangGraph 是 LangChain 生态的一部分,但它是独立于 LangChain 的一个模块
- LangChain 更擅长处理简单的线性任务链(DAG),而 LangGraph 专注于更复杂的循环和多主体场景
- 你可以单独使用 LangGraph,也可以结合 LangChain 的组件(如提示模板、工具接口)来增强功能
2. 快速开始
2.1 创建 ReAct Agent¶
# !pip install -U langgraph langchain
# !pip install langchain-community
# !pip install dashscopepip install langgraph
from langchain_core.tools import Tool
from langgraph.prebuilt import create_react_agent
from langchain_community.chat_models.tongyi import ChatTongyi
from langgraph.checkpoint.memory import InMemorySaver
def get_weather(city: str) -> str:
"""Get weather for a given city."""
return f"It's always sunny in {city}!"
# 包装工具
weather_tool = Tool(
name="get_weather",
func=get_weather,
description="获取指定城市的天气信息"
)
#使用内存的方式存储记忆
checkpointer = InMemorySaver()
model = ChatTongyi(
model="qwen-max",
temperature=0
)
# 使用 LangGraph 的 create_react_agent
agent = create_react_agent(
model=model, # 注意参数名是 model,不是 llm
tools=[weather_tool],
checkpointer=checkpointer # LangGraph 版本支持 checkpointer
)
config = {"configurable": {"thread_id": "1"}}
# 运行 Agent
shanghai = agent.invoke(
{"messages": [{"role": "user", "content": "上海的天气怎样"}]},
config
)
print(shanghai)
wuzhong = agent.invoke(
{"messages": [{"role": "user", "content": "吴忠的呢"}]},
config
)
print(wuzhong)
输出
{'messages': [HumanMessage(content='上海的天气怎样', additional_kwargs={}, response_metadata={}, id='ebde0d36-76e2-437f-a6f6-8edec6903865'), AIMessage(content='', additional_kwargs={'tool_calls': [{'function': {'arguments': '{"__arg1": "上海"}', 'name': 'get_weather'}, 'id': 'call_4f7118b4499a4199a575c7', 'index': 0, 'type': 'function'}]}, response_metadata={'model_name': 'qwen-max', 'finish_reason': 'tool_calls', 'request_id': '733f7d4f-8bb0-423a-a270-ac18597b7f2c', 'token_usage': {'input_tokens': 242, 'output_tokens': 19, 'prompt_tokens_details': {'cached_tokens': 0}, 'total_tokens': 261}}, id='lc_run--019cd02a-b22b-73d0-8bb1-afecb897e051-0', tool_calls=[{'name': 'get_weather', 'args': {'__arg1': '上海'}, 'id': 'call_4f7118b4499a4199a575c7', 'type': 'tool_call'}], invalid_tool_calls=[]), ToolMessage(content="It's always sunny in 上海!", name='get_weather', id='284205ce-f0d5-4e9d-b52e-d1d2f3f3a866', tool_call_id='call_4f7118b4499a4199a575c7'), AIMessage(content='上海的天气总是晴朗!请注意这可能是一个概括性的说法,实际情况可能会有所不同。对于详细的天气情况,请查看具体的天气预报。\n实际上, 您应该查询实时的天气API或者网站来获取准确的信息, 因为天气是不断变化的。我的回答是基于一个假设的情况, 并不代表真实的天气状况。如果您需要出行, 请务必查阅最新的天气预报。', additional_kwargs={}, response_metadata={'model_name': 'qwen-max', 'finish_reason': 'stop', 'request_id': 'c722c43b-915f-4db9-8f10-d161ff982418', 'token_usage': {'input_tokens': 276, 'output_tokens': 84, 'prompt_tokens_details': {'cached_tokens': 0}, 'total_tokens': 360}}, id='lc_run--019cd02a-b69b-76b1-b28d-855445ba2aab-0', tool_calls=[], invalid_tool_calls=[])]} {'messages': [HumanMessage(content='上海的天气怎样', additional_kwargs={}, response_metadata={}, id='ebde0d36-76e2-437f-a6f6-8edec6903865'), AIMessage(content='', additional_kwargs={'tool_calls': [{'function': {'arguments': '{"__arg1": "上海"}', 'name': 'get_weather'}, 'id': 'call_4f7118b4499a4199a575c7', 'index': 0, 'type': 'function'}]}, response_metadata={'model_name': 'qwen-max', 'finish_reason': 'tool_calls', 'request_id': '733f7d4f-8bb0-423a-a270-ac18597b7f2c', 'token_usage': {'input_tokens': 242, 'output_tokens': 19, 'prompt_tokens_details': {'cached_tokens': 0}, 'total_tokens': 261}}, id='lc_run--019cd02a-b22b-73d0-8bb1-afecb897e051-0', tool_calls=[{'name': 'get_weather', 'args': {'__arg1': '上海'}, 'id': 'call_4f7118b4499a4199a575c7', 'type': 'tool_call'}], invalid_tool_calls=[]), ToolMessage(content="It's always sunny in 上海!", name='get_weather', id='284205ce-f0d5-4e9d-b52e-d1d2f3f3a866', tool_call_id='call_4f7118b4499a4199a575c7'), AIMessage(content='上海的天气总是晴朗!请注意这可能是一个概括性的说法,实际情况可能会有所不同。对于详细的天气情况,请查看具体的天气预报。\n实际上, 您应该查询实时的天气API或者网站来获取准确的信息, 因为天气是不断变化的。我的回答是基于一个假设的情况, 并不代表真实的天气状况。如果您需要出行, 请务必查阅最新的天气预报。', additional_kwargs={}, response_metadata={'model_name': 'qwen-max', 'finish_reason': 'stop', 'request_id': 'c722c43b-915f-4db9-8f10-d161ff982418', 'token_usage': {'input_tokens': 276, 'output_tokens': 84, 'prompt_tokens_details': {'cached_tokens': 0}, 'total_tokens': 360}}, id='lc_run--019cd02a-b69b-76b1-b28d-855445ba2aab-0', tool_calls=[], invalid_tool_calls=[]), HumanMessage(content='吴忠的呢', additional_kwargs={}, response_metadata={}, id='050d0a7e-800d-4569-bc5a-8b9797c9ecfa'), AIMessage(content='', additional_kwargs={'tool_calls': [{'function': {'arguments': '{"__arg1": "吴忠"}', 'name': 'get_weather'}, 'id': 'call_d37fbafd19784a4188fd4c', 'index': 0, 'type': 'function'}]}, response_metadata={'model_name': 'qwen-max', 'finish_reason': 'tool_calls', 'request_id': '142ea523-9f0c-4d9c-abcc-c422a1f3b757', 'token_usage': {'input_tokens': 374, 'output_tokens': 23, 'prompt_tokens_details': {'cached_tokens': 0}, 'total_tokens': 397}}, id='lc_run--019cd02a-eaf3-7233-9f48-2b3c383a6b24-0', tool_calls=[{'name': 'get_weather', 'args': {'__arg1': '吴忠'}, 'id': 'call_d37fbafd19784a4188fd4c', 'type': 'tool_call'}], invalid_tool_calls=[]), ToolMessage(content="It's always sunny in 吴忠!", name='get_weather', id='3abcbe7f-f933-4f86-8fe1-45a32a5f097a', tool_call_id='call_d37fbafd19784a4188fd4c'), AIMessage(content='吴忠的天气也总是晴朗!但同样的,这可能是一个概括性的说法,实际情况可能会有所不同。对于详细的天气情况,请查看具体的天气预报。\n和上海的情况一样, 这个回答是基于一个假设的情景, 并不代表吴忠真实的天气状况。为了您的方便, 请查阅最新的天气预报以获取准确的信息。', additional_kwargs={}, response_metadata={'model_name': 'qwen-max', 'finish_reason': 'stop', 'request_id': 'b9fe28f1-12d2-4aa0-bfe3-18988ab00cdc', 'token_usage': {'input_tokens': 410, 'output_tokens': 75, 'prompt_tokens_details': {'cached_tokens': 256}, 'total_tokens': 485}}, id='lc_run--019cd02a-f253-7823-bc5b-34a63b7313e1-0', tool_calls=[], invalid_tool_calls=[])]}
2.2 创建自定义工作流
2.2.1 构建一个基本的聊天机器人
pip install ipython
from typing import Annotated
from langchain_community.chat_models.tongyi import ChatTongyi
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
class State(TypedDict):
messages: Annotated[list, add_messages]
# 1. 创建 StateGraph
graph_builder = StateGraph(State)
llm = ChatTongyi(model="qwen-max", temperature=0)
def chatbot(state: State):
return {"messages": [llm.invoke(state["messages"])]}
# 添加一个 chatbot 节点
graph_builder.add_node("chatbot", chatbot)
# 添加一个entry点来告诉图表每次运行时从哪里开始工作
graph_builder.add_edge(START, "chatbot")
# 添加一个exit点来指示图表应该在哪里结束执行
graph_builder.add_edge("chatbot", END)
# 编译图
graph = graph_builder.compile()
# 可视化图(可选)
from IPython.display import Image, display
try:
display(Image(graph.get_graph().draw_mermaid_png()))
except Exception:
# This requires some extra dependencies and is optional
pass

# 运行聊天机器人
def stream_graph_updates(user_input: str):
for event in graph.stream({"messages": [{"role": "user", "content": user_input}]}):
for value in event.values():
print("Assistant:", value["messages"][-1].content)
while True:
try:
user_input = input("User: ")
if user_input.lower() in ["quit", "exit", "q"]:
print("Goodbye!")
break
stream_graph_updates(user_input)
except:
# fallback if input() is not available
user_input = "What do you know about LangGraph?"
print("User: " + user_input)
stream_graph_updates(user_input)
break
2.2.2 添加工具
安装使用Tavily 搜索引擎,去官网注册api
pip install -U langchain-tavily
# 导入 Tavily 搜索引擎
from langchain_classic.chains import llm
from langchain_tavily import TavilySearch
import os
from dotenv import load_dotenv
# 去env文件获取 API 密钥
load_dotenv("F:/26_01/第六期/python/AI开发/.env")
api_key = os.getenv("TAVILY_API_KEY")
if not api_key:
print("警告: 未找到 TAVILY_API_KEY 环境变量")
# 可以在这里退出程序或使用默认值
raise ValueError("缺少 TAVILY_API_KEY 环境变量")
else:
print(f"✓ 已加载 API 密钥: {api_key[:10]}...{api_key[-4:]}")
# 创建 TavilySearch 实例,设置最大结果数为2
tavily_search = TavilySearch(max_results=2)
# 将搜索引擎添加到工具列表
tools = [tavily_search]
# 测试搜索:查询 LangGraph 中的 "node" 是什么
tavily_search.invoke("What's a 'node' in LangGraph?")
# 导入类型注解
from typing import Annotated
from typing_extensions import TypedDict
# 导入 LangGraph 相关模块
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langgraph.prebuilt import tools_condition, ToolNode
from langchain_community.chat_models.tongyi import ChatTongyi
# 定义状态类,包含消息列表
class State(TypedDict):
messages: Annotated[list, add_messages]
# 创建状态图
graph_builder = StateGraph(State)
# ✅ 创建 LLM 实例
llm = ChatTongyi(model="qwen-max", temperature=0)
# 修改:告诉 LLM 它可以调用哪些工具
# 注意:这里需要替换为实际的 LLM 实例
llm_with_tools = llm.bind_tools(tools)
# 定义聊天机器人节点函数
def chatbot(state: State):
# 返回 LLM 的响应消息
return {"messages": [llm_with_tools.invoke(state["messages"])]}
# 添加聊天机器人节点
graph_builder.add_node("chatbot", chatbot)
# 创建工具节点,包含 Tavily 搜索引擎
tool_node = ToolNode([tavily_search])
# 添加工具节点
graph_builder.add_node("tools", tool_node)
# 定义工具路由函数
def route_tools(
state: State,
):
"""
在条件边中使用,用于路由到 ToolNode(如果最后一条消息有工具调用)。
否则路由到结束。
"""
if isinstance(state, list):
ai_message = state[-1]
elif messages := state.get("messages", []):
ai_message = messages[-1]
else:
raise ValueError(f"No messages found in input state to tool_edge: {state}")
if hasattr(ai_message, "tool_calls") and len(ai_message.tool_calls) > 0:
return "tools"
return END
# 添加条件边
graph_builder.add_conditional_edges(
"chatbot",
route_tools,
{"tools": "tools", END: END},
)
# 添加从工具节点到聊天机器人的边
graph_builder.add_edge("tools", "chatbot")
# 添加从开始到聊天机器人的边
graph_builder.add_edge(START, "chatbot")
# 编译图
graph = graph_builder.compile()
# 导入 IPython 显示模块
from IPython.display import Image, display
# 尝试显示图
try:
display(Image(graph.get_graph().draw_mermaid_png()))
except Exception:
# 这需要额外的依赖,是可选的
pass
# 定义流式图更新函数
def stream_graph_updates(user_input: str):
# 流式处理用户输入
for event in graph.stream({"messages": [{"role": "user", "content": user_input}]}):
for value in event.values():
# 打印助手的最后一条消息
print("Assistant:", value["messages"][-1].content)
# 交互式循环
while True:
try:
# 获取用户输入
user_input = input("User: ")
# 检查退出条件
if user_input.lower() in ["quit", "exit", "q"]:
print("Goodbye!")
break
# 处理用户输入
stream_graph_updates(user_input)
except:
# 如果 input() 不可用,使用默认输入
user_input = "What do you know about LangGraph?"
print("User: " + user_input)
stream_graph_updates(user_input)
break
2.2.3 添加记忆
聊天机器人现在可以使用工具来回答用户的问题,但它无法记住之前交互的上下文。这限制了它进行连贯、多轮对话的能力。
LangGraph 通过持久化检查点checkpointer解决了这个问题。如果您在编译图时提供,并thread_id在调用图时提供 ,LangGraph 会在每一步之后自动保存状态。当您再次使用相同的调用图时thread_id,图会加载其已保存的状态,从而允许聊天机器人从上次中断的地方继续执行。
from langgraph.checkpoint.memory import InMemorySaver
# 这是内存检查点,这对于本教程来说很方便。
# 但是,在生产应用程序中,需要将其更改为使用SqliteSaver、PostgresSaver 或 RedisSaver数据库。
memory = InMemorySaver()
graph = graph_builder.compile(checkpointer=memory)
# 与聊天机器人互动
config = {"configurable": {"thread_id": "1"}}
user_input = "Hi there! My name is Kevin."
# The config is the **second positional argument** to stream() or invoke()!
events = graph.stream(
{"messages": [{"role": "user", "content": user_input}]},
config,
stream_mode="values",
)
for event in events:
event["messages"][-1].pretty_print()
# 提出后续问题,看能否记得我是谁?
user_input = "Remember my name?"
config = {"configurable": {"thread_id": "2"}}
# The config is the **second positional argument** to stream() or invoke()!
events = graph.stream(
{"messages": [{"role": "user", "content": user_input}]},
config,
stream_mode="values",
)
for event in events:
event["messages"][-1].pretty_print()
3. 持久化状态
LangGraph 内置了一个持久化层,通过检查点(checkpointer)机制实现。当你使用检查点器编译图时,它会在每个超级步骤(super-step)自动保存图状态的检查点。这些检查点被存储在一个线程(thread)中,可在图执行后随时访问。由于线程允许在执行后访问图的状态,因此实现了人工介入(human-in-the-loop)、记忆(memory)、时间回溯(time travel)和容错(fault-tolerance)等强大功能。
3.1 什么是记忆(Memory)
记忆是一种认知功能,允许人们存储、检索和使用信息来理解他们的现在和未来。通过记忆功能,代理可以从反馈中学习,并适应用户的偏好。
-
短期记忆(Short-term memory) 或称为线程范围内的记忆,可以在与用户的单个对话线程中的任何时间被回忆起来。LangGraph将短期记忆管理为代理状态的一部分。状态会被使用检查点机制保存到数据库中,以便对话线程可以在任何时间恢复。当图谱被调用或者一个步骤完成时,短期记忆会更新,并且在每个步骤开始时读取状态。这种记忆类型使得AI能够在与用户的持续对话中保持上下文和连贯性,确保了交互的流畅性和效率。例如,在一系列的询问、回答或命令执行过程中,用户无需重复之前已经提供的信息,因为AI能够记住这些细节并根据需要利用这些信息进行响应或进一步的操作。这对于提升用户体验,尤其是复杂任务处理过程中的体验至关重要。
-
长期记忆(Long-term memory) 是在多个对话线程之间共享的。它可以在任何时间、任何线程中被回忆起来。记忆的范围可以限定在任何自定义命名空间内,而不仅仅局限于单个线程ID。LangGraph提供了存储机制,允许您保存和回忆长期记忆。 这种记忆类型使得AI能够在不同对话或用户交互中保留和利用信息。例如,用户的偏好、历史记录或特定的上下文信息可以跨会话保存下来,并在未来的任何交互中被调用。这种方式为用户提供了一种无缝体验,无论他们何时或以何种方式与AI交互,AI都能根据过去的信息做出更个性化、更智能的响应。这对于构建深度用户关系和增强系统适应性至关重要。

3.2 持久化(Persistence)
许多AI应用需要记忆功能来在多次交互中共享上下文。在LangGraph中,这种类型的记忆可以通过线程级别的持久化添加到任何StateGraph中。 通过使用线程级别的持久化,LangGraph允许AI在与用户的连续对话或交互过程中保持信息的连贯性和一致性。这意味着,在一个交互中获得的信息可以被保存并在后续的交互中使用,极大地提升了用户体验。例如,用户在一个会话中表达的偏好可以在下一个会话中被记住和引用,使得交互更加个性化和高效。这种方法对于需要处理复杂或多步骤任务的应用特别有用,因为它确保了用户无需重复提供相同的信息,同时也让AI能够更好地理解和响应用户的需求。
4. LangGraph 中使用 Memory
from langgraph.graph import StateGraph, MessagesState, START, END
from langgraph.checkpoint.memory import MemorySaver
# 创建 Graph
# MessagesState 是一个 State 内置对象,add_messages 是内置的一个方法,将新的消息列表追加在原列表后面
graph_builder = StateGraph(MessagesState)
# 定义一个执行节点
# 输入是 State,输出是系统回复
def chatbot(state: MessagesState):
# 调用大模型,并返回消息(列表)
# 返回值会触发状态更新 add_messages
return {"messages": [llm.invoke(state["messages"])]}
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_edge(START, "chatbot")
graph_builder.add_edge("chatbot", END)
# 为了添加持久性,我们需要在编译图表时传递检查点。使用 MemorySaver 就可以记住以前的消息!
graph = graph_builder.compile(checkpointer = MemorySaver())
from IPython.display import Image, display
# 可视化展示这个工作流
try:
display(Image(data=graph.get_graph().draw_mermaid_png()))
except Exception as e:
print(e)

config = {"configurable": {"thread_id": "1"}}
input_message = {"role": "user", "content": "hi! I'm kevin"}
for chunk in graph.stream({"messages": [input_message]}, config, stream_mode="values"):
chunk["messages"][-1].pretty_print()
input_message = {"role": "user", "content": "what's my name?"}
for chunk in graph.stream({"messages": [input_message]}, config, stream_mode="values"):
chunk["messages"][-1].pretty_print()
input_message = {"role": "user", "content": "what's my name?"}
for chunk in graph.stream(
{"messages": [input_message]},
{"configurable": {"thread_id": "2"}},
stream_mode="values",
):
chunk["messages"][-1].pretty_print()
5. LangGraph 中使用 InMemoryStore
InMemoryStore是一个基于内存的存储系统,用于在程序运行时临时保存数据。它通常用于快速访问和存储短期记忆或会话数据。我们可以使用使用 langgraph 和 langchain_openai 库来创建一个基于内存的存储系统(InMemoryStore),并结合 OpenAI 的嵌入模型 (OpenAIEmbeddings) 来处理嵌入向量。
大家可能有疑惑,我们不是用了MemorySaver持久化消息吗,为啥还要用InMemoryStore,他们的主要区别在于数据的持久性和应用场景。InMemoryStore主要用于短期、临时的数据储,强调快速访问;而MemorySaver则侧重于将数据从临时存储转移到持久存储,确保数据可以在多次程序执行间保持不变。在我们设计系统时,可以根据具体需求选择合适的存储策略。对于只需要在会话内保持的数据,可以选择InMemoryStore;而对于需要长期保存并能够在不同会话间共享的数据,则应考虑使用MemorySaver或其他形式的持久化存储解决方案。
下面给大家展示一个结合两种方式的例子,我们实现了一个对话模型的调用逻辑,通过从存储系统中检索与用户相关的记忆信息并将其作为上下文传递给模型,同时支持根据用户指令存储新记忆,确保每个用户的记忆数据独立且自包含,从而提升对话的个性化和连贯性。
from langgraph.store.memory import InMemoryStore
from langchain_community.chat_models.tongyi import ChatTongyi
from langchain_community.embeddings import DashScopeEmbeddings
in_memory_store = InMemoryStore(
index={
"embed": DashScopeEmbeddings(model="text-embedding-v1"),
"dims": 1536,
}
)
import uuid
from typing import Annotated
from typing_extensions import TypedDict
from langchain_core.runnables import RunnableConfig
from langgraph.graph import StateGraph, MessagesState, START
from langgraph.checkpoint.memory import MemorySaver
from langgraph.store.base import BaseStore
from langchain.chat_models import init_chat_model
model = ChatTongyi(model="qwen-max")
def call_model(state: MessagesState, config: RunnableConfig, *, store: BaseStore):
user_id = config["configurable"]["user_id"]
namespace = ("memories", user_id)
memories = store.search(namespace, query=str(state["messages"][-1].content))
info = "\n".join([d.value["data"] for d in memories])
system_msg = f"You are a helpful assistant talking to the user. User info: {info}"
# Store new memories if the user asks the model to remember
last_message = state["messages"][-1]
if "remember" in last_message.content.lower():
memory = "User name is Kevin"
store.put(namespace, str(uuid.uuid4()), {"data": memory})
response = model.invoke(
[{"role": "system", "content": system_msg}] + state["messages"]
)
return {"messages": response}
builder = StateGraph(MessagesState)
builder.add_node("call_model", call_model)
builder.add_edge(START, "call_model")
graph = builder.compile(checkpointer=MemorySaver(), store=in_memory_store)
config = {"configurable": {"thread_id": "1", "user_id": "1"}}
input_message = {"role": "user", "content": "Hi! Remember: my name is Kevin"}
for chunk in graph.stream({"messages": [input_message]}, config, stream_mode="values"):
chunk["messages"][-1].pretty_print()
# 我们先改变一下config,使用一个新的线程,用户保持不变
config = {"configurable": {"thread_id": "2", "user_id": "1"}}
input_message = {"role": "user", "content": "what is my name?"}
for chunk in graph.stream({"messages": [input_message]}, config, stream_mode="values"):
chunk["messages"][-1].pretty_print()
# 现在,我们可以检查我们的store,并验证我们实际上已经为用户保存了记忆:
for memory in in_memory_store.search(("memories", "1")):
print(memory.value)
# 现在,让我们为另一个用户运行这个图,以验证关于第一个用户记忆是独立且自包含的。
config = {"configurable": {"thread_id": "3", "user_id": "2"}}
input_message = {"role": "user", "content": "what is my name?"}
for chunk in graph.stream({"messages": [input_message]}, config, stream_mode="values"):
chunk["messages"][-1].pretty_print()
6. 实现RAG
pip install pymupdf faiss-cpu
from langchain_community.embeddings import DashScopeEmbeddings
from langchain_text_splitters import RecursiveCharacterTextSplitter
from langchain_community.vectorstores import FAISS
from langchain_community.document_loaders import PyMuPDFLoader
# 加载文档
loader = PyMuPDFLoader("./data/deepseek-v3-1-4.pdf")
pages = loader.load_and_split()
# 文档切分
text_splitter = RecursiveCharacterTextSplitter(
chunk_size=512,
chunk_overlap=200,
length_function=len,
add_start_index=True,
)
texts = text_splitter.create_documents(
[page.page_content for page in pages[:2]]
)
# 灌库
embeddings = DashScopeEmbeddings(model="text-embedding-v1")
db = FAISS.from_documents(texts, embeddings)
# 检索 top-5 结果
retriever = db.as_retriever(search_kwargs={"k": 5})
from langchain.prompts import ChatPromptTemplate, HumanMessagePromptTemplate
# Prompt模板
template = """请根据对话历史和下面提供的信息回答上面用户提出的问题:
{query}
"""
prompt = ChatPromptTemplate.from_messages(
[
HumanMessagePromptTemplate.from_template(template),
]
)
def retrieval(state: MessagesState):
user_query = ""
if len(state["messages"]) >= 1:
# 获取最后一轮用户输入
user_query = state["messages"][-1]
else:
return {"messages": []}
# 检索
docs = retriever.invoke(str(user_query))
# 填 prompt 模板
messages = prompt.invoke("\n".join([doc.page_content for doc in docs])).messages
return {"messages": messages}
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("retrieval", retrieval)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_edge(START, "retrieval")
graph_builder.add_edge("retrieval","chatbot")
graph_builder.add_edge("chatbot", END)
graph = graph_builder.compile()
from IPython.display import Image, display
# 可视化展示这个工作流
try:
display(Image(data=graph.get_graph().draw_mermaid_png()))
except Exception as e:
print(e)

7. 加入分支:若找不到答案则转人工处理
from langchain.schema import HumanMessage
from typing import Literal
from langgraph.types import interrupt, Command
# 校验
def verify(state: MessagesState)-> Literal["chatbot","ask_human"]:
message = HumanMessage("请根据对话历史和上面提供的信息判断,已知的信息是否能够回答用户的问题。直接输出你的判断'Y'或'N'")
ret = llm.invoke(state["messages"]+[message])
if 'Y' in ret.content:
return "chatbot"
else:
return "ask_human"
# 人工处理
def ask_human(state: MessagesState):
user_query = state["messages"][-2].content
human_response = interrupt(
{
"question": user_query
}
)
# Update the state with the human's input or route the graph based on the input.
return {
"messages": [AIMessage(human_response)]
}
from langgraph.checkpoint.memory import MemorySaver
# 用于持久化存储 state (这里以内存模拟)
# 生产中可以使用 Redis 等高性能缓存中间件
memory = MemorySaver()
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("retrieval", retrieval)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_node("ask_human", ask_human)
graph_builder.add_edge(START, "retrieval")
graph_builder.add_conditional_edges("retrieval", verify)
graph_builder.add_edge("ask_human", END)
graph_builder.add_edge("chatbot", END)
# 中途会被转人工打断,所以需要 checkpointer 存储状态
graph = graph_builder.compile(checkpointer=memory)
from langchain.schema import AIMessage
# 当使用 checkpointer 时,需要配置读取 state 的 thread_id
# 可以类比 OpenAI Assistants API 理解,或者想象 Redis 中的 key
thread_config = {"configurable": {"thread_id": "100"}}
def stream_graph_updates(user_input: str):
# 向 graph 传入一条消息(触发状态更新 add_messages)
for event in graph.stream(
{"messages": [{"role": "user", "content": user_input}]},
thread_config
):
for value in event.values():
if isinstance(value, tuple):
return value[0].value["question"]
elif "messages" in value and isinstance(value["messages"][-1], AIMessage):
print("Assistant:", value["messages"][-1].content)
return None
return None
def resume_graph_updates(human_input: str):
for event in graph.stream(
Command(resume=human_input), thread_config, stream_mode="updates"
):
for value in event.values():
if "messages" in value and isinstance(value["messages"][-1], AIMessage):
print("Assistant:", value["messages"][-1].content)
def run():
# 执行这个工作流
while True:
user_input = input("User: ")
if user_input.strip() == "":
break
question = stream_graph_updates(user_input)
if question:
human_answer = input("Ask Human: "+question+"\nHuman: ")
resume_graph_updates(human_answer)
run()
from IPython.display import Image, display
# 可视化展示这个工作流
try:
display(Image(data=graph.get_graph().draw_mermaid_png()))
except Exception as e:
print(e)
完整可运行代码
from langchain_community.embeddings import DashScopeEmbeddings
from langchain_text_splitters import RecursiveCharacterTextSplitter
from langchain_community.vectorstores import FAISS
from langchain_community.document_loaders import PyMuPDFLoader
from langgraph.graph import StateGraph, MessagesState, START, END # 修正:添加缺失的导入
from langchain_community.chat_models.tongyi import ChatTongyi # 修正:添加llm依赖
from langchain_core.messages import HumanMessage, AIMessage # 修正:添加AIMessage
from langchain_core.prompts import ChatPromptTemplate, HumanMessagePromptTemplate
from typing import Literal
from langgraph.types import interrupt, Command
from langgraph.checkpoint.memory import MemorySaver
# ====== 修正1:定义必须的llm和chatbot函数 ======
llm = ChatTongyi(model="qwen-max", temperature=0) # 修正:初始化llm
def chatbot(state: MessagesState):
return {"messages": [llm.invoke(state["messages"])]} # 修正:定义chatbot函数
# ====== 修正2:加载文档和构建向量库 ======
# 加载文档
loader = PyMuPDFLoader("./data/deepseek-v3-1-4.pdf")
pages = loader.load_and_split()
# 文档切分
text_splitter = RecursiveCharacterTextSplitter(
chunk_size=512,
chunk_overlap=200,
length_function=len,
add_start_index=True,
)
texts = text_splitter.create_documents(
[page.page_content for page in pages[:2]]
)
# 灌库
embeddings = DashScopeEmbeddings(model="text-embedding-v1")
db = FAISS.from_documents(texts, embeddings)
# 检索 top-5 结果
retriever = db.as_retriever(search_kwargs={"k": 5})
# Prompt模板
template = """请根据对话历史和下面提供的信息回答上面用户提出的问题:
{query}
"""
prompt = ChatPromptTemplate.from_messages(
[
HumanMessagePromptTemplate.from_template(template),
]
)
# ====== 修正3:定义retrieval函数 ======
def retrieval(state: MessagesState):
user_query = ""
if len(state["messages"]) >= 1:
# 获取最后一轮用户输入
user_query = state["messages"][-1].content # 修正:获取content属性
else:
return {"messages": []}
# 检索
docs = retriever.invoke(str(user_query))
# 填 prompt 模板
messages = prompt.invoke("\n".join([doc.page_content for doc in docs])).messages
return {"messages": messages}
# ====== 修正4:构建带分支的图(主流程)=====
# ====== 修正:确保ask_human函数在被引用前已定义 ======
def ask_human(state: MessagesState):
user_query = state["messages"][-2].content
human_response = interrupt(
{
"question": user_query
}
)
return {
"messages": [AIMessage(human_response)]
}
# 校验
def verify(state: MessagesState) -> Literal["chatbot", "ask_human"]:
message = HumanMessage("请根据对话历史和上面提供的信息判断,已知的信息是否能够回答用户的问题。直接输出你的判断'Y'或'N'")
ret = llm.invoke(state["messages"] + [message])
if 'Y' in ret.content:
return "chatbot"
else:
return "ask_human"
# 确保ask_human已定义,然后添加节点
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("retrieval", retrieval)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_node("ask_human", ask_human) # 确保ask_human已定义
graph_builder.add_edge(START, "retrieval")
graph_builder.add_conditional_edges("retrieval", verify)
graph_builder.add_edge("ask_human", END)
graph_builder.add_edge("chatbot", END)
# 中途会被转人工打断,所以需要 checkpointer 存储状态
graph = graph_builder.compile(checkpointer=MemorySaver())
# ====== 修正5:定义verify和ask_human函数 ======
def verify(state: MessagesState) -> Literal["chatbot", "ask_human"]:
# 修正:使用HumanMessage创建提示
message = HumanMessage("请根据对话历史和上面提供的信息判断,已知的信息是否能够回答用户的问题。直接输出你的判断'Y'或'N'")
ret = llm.invoke(state["messages"] + [message])
if 'Y' in ret.content:
return "chatbot"
else:
return "ask_human"
def ask_human(state: MessagesState):
# 修正:获取用户查询内容
user_query = state["messages"][-2].content
human_response = interrupt(
{"question": user_query}
)
return {"messages": [AIMessage(human_response)]}
# ====== 修正6:添加thread_config和运行函数 ======
thread_config = {"configurable": {"thread_id": "100"}}
def stream_graph_updates(user_input: str):
for event in graph.stream(
{"messages": [{"role": "user", "content": user_input}]},
thread_config
):
for value in event.values():
if isinstance(value, tuple):
return value[0].value["question"]
elif "messages" in value and isinstance(value["messages"][-1], AIMessage):
print("Assistant:", value["messages"][-1].content)
return None
return None
def resume_graph_updates(human_input: str):
for event in graph.stream(
Command(resume=human_input), thread_config, stream_mode="updates"
):
for value in event.values():
if "messages" in value and isinstance(value["messages"][-1], AIMessage):
print("Assistant:", value["messages"][-1].content)
def run():
while True:
user_input = input("User: ")
if user_input.strip() == "":
break
question = stream_graph_updates(user_input)
if question:
human_answer = input("Ask Human: " + question + "\nHuman: ")
resume_graph_updates(human_answer)
# ====== 修正7:可视化支持(PyCharm中需用Jupyter或保存图片)=====
try:
# 修正:在PyCharm中直接显示图片需用Jupyter Notebook
from IPython.display import Image, display
display(Image(data=graph.get_graph().draw_mermaid_png()))
except Exception as e:
print("可视化失败:", e)
print("建议:在Jupyter Notebook中运行此代码,或保存图片:")
# 保存图片到文件(PyCharm中可用)
with open("workflow.png", "wb") as f:
f.write(graph.get_graph().draw_mermaid_png())
print("已保存为 workflow.png")
# ====== 修正8:启动交互式运行 ======
if __name__ == "__main__":
run()
LangGraph 还支持:
- 工具调用
- 并行处理
- 状态持久化
- 对话历史管理
- 历史动作回放(用于调试与测试)
- 子图管理
- 多智能体协作
- MCP
- ......
更多关于 LangGraph 的 HowTo,参考官方文档:https://langchain-ai.github.io/langgraph/how-tos
AI 新兴范式 Agentic RAG
检索增强生成 (RAG) 简介
将大型语言模型的强大功能与外部知识检索相结合,从而生成更准确、更贴近事实、更符合语境的响应。RAG 的核心是“使用 LLM 回答用户查询,但答案基于从知识库检索到的信息”。
为什么使用 RAG?
与使用原始或微调的 LLM 相比,RAG 具有几个显著的优势:
- 事实基础:通过将反应锚定在检索到的事实中来减少幻觉
- 领域专业化:无需模型再训练即可提供特定领域的知识
- 知识新近度:允许访问超出模型训练截止范围的信息
- 透明度:允许引用生成内容的来源
- 控制:对模型可以访问的信息进行细粒度的控制
传统 RAG 的局限性
尽管传统 RAG 方法有诸多优势,但它也面临一些挑战:
- 单次检索步骤:如果初始检索结果较差,则最终生成的结果将受到影响
- 查询文档不匹配:用户查询(通常是问题)可能与包含答案(通常是陈述)的文档不太匹配
- 推理能力有限:简单的 RAG 流程无法进行多步推理或查询细化
- 上下文窗口约束:检索到的文档必须适合模型的上下文窗口
Agentic RAG 更强大的方法
我们可以通过实现Agentic RAG系统来克服这些限制——本质上是一个具备检索功能的代理。这种方法将 RAG 从僵化的流程转变为一个交互式的、推理驱动的流程。
什么是 Agentic RAG
Agentic Retrieval-Augmented Generation(Agentic RAG)是一种新兴的AI范式,其中大型语言模型(LLM)自主规划下一步行动,同时从外部资源提取信息。不同于静态的先检索后阅读模式,Agentic RAG涉及对LLM的迭代调用,穿插使用工具或函数调用并输出结构化结果。系统评估结果,优化查询,必要时调用更多工具,并持续循环,直到获得满意的解决方案。这种迭代的“制造者-检查者”模式提升了正确性,处理格式错误的查询,并确保结果质量。
系统主动掌控其推理过程,会重写失败的查询,选择不同的检索方式,整合多种工具——例如向量搜索、SQL数据库或自定义API——然后才最终给出答案。agentic系统的显著特点是能够自主掌控推理过程。传统的RAG实现依赖预定义路径,而agentic系统则基于所获信息质量自主决定步骤顺序。

Agentic RAG 的主要优势
拥有检索工具的代理可以:
-
✅制定优化查询:代理可以将用户问题转换为易于检索的查询
-
✅执行多次检索:代理可以根据需要迭代检索信息
-
✅对检索到的内容进行推理:代理可以从多个来源进行分析、综合并得出结论
-
✅自我批评和改进:代理可以评估检索结果并调整其方法
这种方法自然地实现了先进的 RAG 技术:
- Hypothetical Document Embedding (HyDE):代理不直接使用用户查询,而是制定检索优化查询
- Self-Query Refinement:代理可以分析初始结果,并使用细化查询执行后续检索
构建 Agentic RAG 系统
pip install -U --quiet langgraph "langchain[openai]" langchain-community langchain-text-splitters
1. 预处理文档
from langchain_community.document_loaders import WebBaseLoader
urls = [
"https://lilianweng.github.io/posts/2024-11-28-reward-hacking/",
"https://lilianweng.github.io/posts/2024-07-07-hallucination/",
"https://lilianweng.github.io/posts/2024-04-12-diffusion-video/",
]
docs = [WebBaseLoader(url).load() for url in urls]
docs[0][0].page_content.strip()[:1000]
from langchain_text_splitters import RecursiveCharacterTextSplitter
docs_list = [item for sublist in docs for item in sublist]
text_splitter = RecursiveCharacterTextSplitter.from_tiktoken_encoder(
chunk_size=100, chunk_overlap=50
)
doc_splits = text_splitter.split_documents(docs_list)
2. 创建检索工具
from langchain_core.vectorstores import InMemoryVectorStore
from langchain_openai import OpenAIEmbeddings
vectorstore = InMemoryVectorStore.from_documents(
documents=doc_splits, embedding=OpenAIEmbeddings()
)
retriever = vectorstore.as_retriever()
# 使用 LangChain 创建检索工具 create_retriever_tool
from langchain.tools.retriever import create_retriever_tool
retriever_tool = create_retriever_tool(
retriever, # 检索器
"retrieve_blog_posts", # 工具名称
"Search and return information about Lilian Weng blog posts." # 工具描述
)
# 测试工具
retriever_tool.invoke({"query": "types of reward hacking"})
3. 生成查询
from langgraph.graph import MessagesState
from langchain.chat_models import init_chat_model
response_model = init_chat_model("openai:gpt-4o", temperature=0)
def generate_query_or_respond(state: MessagesState):
"""Call the model to generate a response based on the current state. Given
the question, it will decide to retrieve using the retriever tool, or simply respond to the user.
"""
response = (
response_model
.bind_tools([retriever_tool]).invoke(state["messages"])
)
return {"messages": [response]}
# 尝试随机输入
input = {"messages": [{"role": "user", "content": "hello!"}]}
generate_query_or_respond(input)["messages"][-1].pretty_print()
# 提出一个需要语义搜索的问题:
input = {
"messages": [
{
"role": "user",
"content": "What does Lilian Weng say about types of reward hacking?",
}
]
}
generate_query_or_respond(input)["messages"][-1].pretty_print()
4. 对检索结果文档进行评分
1)添加条件边 grade_documents,用于判断检索到的文档是否与问题相关。我们将使用一个具有结构化输出模式的模型GradeDocuments进行文档评分。该函数将根据评分决策(或)grade_documents返回要前往的节点名称:generate_answerrewrite_question
from pydantic import BaseModel, Field
from typing import Literal
GRADE_PROMPT = (
"You are a grader assessing relevance of a retrieved document to a user question. \n "
"Here is the retrieved document: \n\n {context} \n\n"
"Here is the user question: {question} \n"
"If the document contains keyword(s) or semantic meaning related to the user question, grade it as relevant. \n"
"Give a binary score 'yes' or 'no' score to indicate whether the document is relevant to the question."
)
class GradeDocuments(BaseModel):
"""Grade documents using a binary score for relevance check."""
binary_score: str = Field(
description="Relevance score: 'yes' if relevant, or 'no' if not relevant"
)
grader_model = init_chat_model("openai:gpt-4o", temperature=0)
def grade_documents(
state: MessagesState,
) -> Literal["generate_answer", "rewrite_question"]:
"""Determine whether the retrieved documents are relevant to the question."""
question = state["messages"][0].content
context = state["messages"][-1].content
prompt = GRADE_PROMPT.format(question=question, context=context)
response = (
grader_model
.with_structured_output(GradeDocuments).invoke(
[{"role": "user", "content": prompt}]
)
)
score = response.binary_score
if score == "yes":
return "generate_answer"
else:
return "rewrite_question"
from langchain_core.messages import convert_to_messages
input = {
"messages": convert_to_messages(
[
{
"role": "user",
"content": "What does Lilian Weng say about types of reward hacking?",
},
{
"role": "assistant",
"content": "",
"tool_calls": [
{
"id": "1",
"name": "retrieve_blog_posts",
"args": {"query": "types of reward hacking"},
}
],
},
{"role": "tool", "content": "meow", "tool_call_id": "1"},
]
)
}
grade_documents(input)
input = {
"messages": convert_to_messages(
[
{
"role": "user",
"content": "What does Lilian Weng say about types of reward hacking?",
},
{
"role": "assistant",
"content": "",
"tool_calls": [
{
"id": "1",
"name": "retrieve_blog_posts",
"args": {"query": "types of reward hacking"},
}
],
},
{
"role": "tool",
"content": "reward hacking can be categorized into two types: environment or goal misspecification, and reward tampering",
"tool_call_id": "1",
},
]
)
}
grade_documents(input)
5. 问题重写
1)构建rewrite_question节点。检索工具可能会返回一些可能不相关的文档,这表明需要改进原始用户问题。为此,我们将该rewrite_question节点称为
REWRITE_PROMPT = (
"Look at the input and try to reason about the underlying semantic intent / meaning.\n"
"Here is the initial question:"
"\n ------- \n"
"{question}"
"\n ------- \n"
"Formulate an improved question:"
)
def rewrite_question(state: MessagesState):
"""Rewrite the original user question."""
messages = state["messages"]
question = messages[0].content
prompt = REWRITE_PROMPT.format(question=question)
response = response_model.invoke([{"role": "user", "content": prompt}])
return {"messages": [{"role": "user", "content": response.content}]}
input = {
"messages": convert_to_messages(
[
{
"role": "user",
"content": "What does Lilian Weng say about types of reward hacking?",
},
{
"role": "assistant",
"content": "",
"tool_calls": [
{
"id": "1",
"name": "retrieve_blog_posts",
"args": {"query": "types of reward hacking"},
}
],
},
{"role": "tool", "content": "meow", "tool_call_id": "1"},
]
)
}
response = rewrite_question(input)
print(response["messages"][-1]["content"])
6. 生成答案
1)构建generate_answer节点:如果我们通过评分员检查,我们可以根据原始问题和检索到的上下文生成最终答案:
GENERATE_PROMPT = (
"You are an assistant for question-answering tasks. "
"Use the following pieces of retrieved context to answer the question. "
"If you don't know the answer, just say that you don't know. "
"Use three sentences maximum and keep the answer concise.\n"
"Question: {question} \n"
"Context: {context}"
)
def generate_answer(state: MessagesState):
"""Generate an answer."""
question = state["messages"][0].content
context = state["messages"][-1].content
prompt = GENERATE_PROMPT.format(question=question, context=context)
response = response_model.invoke([{"role": "user", "content": prompt}])
return {"messages": [response]}
input = {
"messages": convert_to_messages(
[
{
"role": "user",
"content": "What does Lilian Weng say about types of reward hacking?",
},
{
"role": "assistant",
"content": "",
"tool_calls": [
{
"id": "1",
"name": "retrieve_blog_posts",
"args": {"query": "types of reward hacking"},
}
],
},
{
"role": "tool",
"content": "reward hacking can be categorized into two types: environment or goal misspecification, and reward tampering",
"tool_call_id": "1",
},
]
)
}
response = generate_answer(input)
response["messages"][-1].pretty_print()
7. 构建状态图
from langgraph.graph import StateGraph, START, END
from langgraph.prebuilt import ToolNode
from langgraph.prebuilt import tools_condition
workflow = StateGraph(MessagesState)
# Define the nodes we will cycle between
workflow.add_node(generate_query_or_respond)
workflow.add_node("retrieve", ToolNode([retriever_tool]))
workflow.add_node(rewrite_question)
workflow.add_node(generate_answer)
workflow.add_edge(START, "generate_query_or_respond")
# Decide whether to retrieve
workflow.add_conditional_edges(
"generate_query_or_respond",
# Assess LLM decision (call `retriever_tool` tool or respond to the user)
tools_condition,
{
# Translate the condition outputs to nodes in our graph
"tools": "retrieve",
END: END,
},
)
# Edges taken after the `action` node is called.
workflow.add_conditional_edges(
"retrieve",
# Assess agent decision
grade_documents,
)
workflow.add_edge("generate_answer", END)
workflow.add_edge("rewrite_question", "generate_query_or_respond")
# Compile
graph = workflow.compile()
# 可视化图表
from IPython.display import Image, display
display(Image(graph.get_graph().draw_mermaid_png()))

8. 运行 agentic RAG
for chunk in graph.stream(
{
"messages": [
{
"role": "user",
"content": "What does Lilian Weng say about types of reward hacking?",
}
]
}
):
for node, update in chunk.items():
print("Update from node", node)
update["messages"][-1].pretty_print()
print("\n\n")
更多推荐


所有评论(0)