AI Agent工作流编排是当前大模型开发的核心能力之一。LangGraph作为LangChain团队推出的状态图编排框架,通过有向无环图(DAG)管理多轮对话流程,结合Function Calling实现工具调用链路,构建出可循环、可分支的智能体工作流。本文围绕LangGraph状态图模型、工具注册调用、节点路由和记忆持久化四个环节,给出完整的代码实现。
LangGraph状态图模型与多轮对话编排机制
LangGraph的核心抽象是StateGraph。每个节点接收当前状态(State),执行逻辑后返回状态更新。状态使用TypedDict定义,所有节点共享同一状态结构。
from typing import Annotated, TypedDict
from langgraph.graph import StateGraph, END
from langchain_core.messages import BaseMessage
import operator
class AgentState(TypedDict):
messages: Annotated[list[BaseMessage], operator.add]
tool_results: list[str]
iteration: int
def create_graph():
workflow = StateGraph(AgentState)
# 注册节点
workflow.add_node("planner", plan_node)
workflow.add_node("executor", execute_node)
workflow.add_node("reviewer", review_node)
# 设置入口
workflow.set_entry_point("planner")
# 条件路由:planner -> executor 或 END
workflow.add_conditional_edges(
"planner",
lambda state: "executor" if state["iteration"] < 3 else END,
{"executor": "executor", END: END}
)
# executor -> reviewer
workflow.add_edge("executor", "reviewer")
# reviewer -> planner(循环)或 END
workflow.add_conditional_edges(
"reviewer",
lambda state: "planner" if not state["review_passed"] else END,
{"planner": "planner", END: END}
)
return workflow.compile()
planner节点负责拆解任务,executor节点调用工具执行,reviewer节点验证结果是否达标。三节点循环结构覆盖了”规划-执行-验证”的完整Agent闭环。
Function Calling工具注册与调用链路实现
工具定义使用@tool装饰器,LangGraph在节点内部自动处理Function Calling的请求-响应流程。
from langchain_core.tools import tool
from langchain_openai import ChatOpenAI
@tool
def search_database(query: str, limit: int = 10) -> str:
'''搜索业务数据库,返回匹配的记录。'''
results = db.search(query, limit=limit)
return json.dumps(results, ensure_ascii=False)
@tool
def calculate_metrics(data: str) -> str:
'''对给定数据计算统计指标。'''
values = json.loads(data)
return json.dumps({
"mean": sum(values) / len(values),
"max": max(values),
"min": min(values)
})
# 绑定工具到LLM
llm = ChatOpenAI(model="gpt-4o", temperature=0)
llm_with_tools = llm.bind_tools([search_database, calculate_metrics])
def execute_node(state: AgentState) -> dict:
response = llm_with_tools.invoke(state["messages"])
tool_results = []
if response.tool_calls:
for tc in response.tool_calls:
if tc["name"] == "search_database":
result = search_database.invoke(tc["args"])
elif tc["name"] == "calculate_metrics":
result = calculate_metrics.invoke(tc["args"])
tool_results.append(result)
from langchain_core.messages import ToolMessage
state["messages"].append(
ToolMessage(content=result, tool_call_id=tc["id"])
)
return {
"messages": [response],
"tool_results": tool_results,
"iteration": state["iteration"] + 1
}
工具调用的关键在于ToolMessage的构造。每个tool_call携带唯一的tool_call_id,对应的ToolMessage必须使用相同的id,LLM才能正确关联调用与结果。
LangGraph节点编排与条件路由配置
实际业务中,Agent需要根据中间结果动态选择执行路径。条件边(conditional_edges)通过回调函数返回的字符串值决定下一跳节点。
def review_node(state: AgentState) -> dict:
'''验证执行结果是否满足要求'''
last_result = state["tool_results"][-1] if state["tool_results"] else ""
data = json.loads(last_result) if last_result else {}
review_passed = False
if data.get("max", 0) > 100:
review_passed = True
from langchain_core.messages import AIMessage
return {
"messages": [AIMessage(content=f"Review: {'passed' if review_passed else 'failed'}")],
"review_passed": review_passed
}
# 编译时可加入中断点,支持人工审批
graph = create_graph()
config = {"configurable": {"thread_id": "session-001"}}
result = graph.invoke(
{"messages": [HumanMessage(content="查询销售数据并计算指标")], "iteration": 0},
config=config
)
interrupt_before参数让图在指定节点前暂停,配合checkpointer实现人机协作。这在需要人工审批的高风险场景(如资金操作、数据删除)中必不可少。
Agent记忆持久化与上下文窗口管理
长时间运行的Agent会积累大量对话历史,超出上下文窗口限制。通过MemorySaver检查点和消息裁剪策略管理上下文。
from langgraph.checkpoint.memory import MemorySaver
from langchain_core.messages import trim_messages
checkpointer = MemorySaver()
def create_graph_with_memory():
workflow = StateGraph(AgentState)
workflow.add_node("planner", plan_node)
workflow.add_node("executor", execute_node)
workflow.add_node("reviewer", review_node)
workflow.set_entry_point("planner")
workflow.add_conditional_edges("planner", route_after_plan)
workflow.add_edge("executor", "reviewer")
workflow.add_conditional_edges("reviewer", route_after_review)
return workflow.compile(checkpointer=checkpointer)
# 消息裁剪策略:保留最近20条消息
def trim_messages_if_needed(state: AgentState) -> dict:
trimmer = trim_messages(
max_tokens=4000,
strategy="last",
token_counter=len,
start_on="human",
end_on="ai"
)
trimmed = trimmer.invoke(state["messages"])
return {"messages": trimmed}
# 跨会话复用:通过thread_id区分不同会话
config_user_a = {"configurable": {"thread_id": "user-a-session"}}
graph.invoke(
{"messages": [HumanMessage(content="查询上月数据")]},
config=config_user_a
)
MemorySaver将状态存储在内存中,生产环境可替换为PostgresSaver或RedisSaver实现持久化。配合消息裁剪,Agent可以在长对话中保持稳定的响应质量,避免因上下文溢出导致工具调用失败。
以上四个环节构成了LangGraph Agent工作流的完整链路。状态图定义流程骨架,Function Calling实现工具调用,条件路由支持动态决策,记忆管理保障长期运行稳定性。实际部署时还需加入超时控制、重试机制和调用链路追踪,确保Agent在生产环境中的可靠性。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/aiagent-gong-zuo-liu-bian-pai-shi-zhan-langgraph-zhuang-tai/