大模型Agent工作流编排是AI应用落地的关键环节。LangGraph作为LangChain生态中专注于有状态多步骤工作流的框架,提供了循环执行、条件分支、状态持久化等能力,适合构建需要多轮推理和工具调用的复杂Agent系统。本文围绕LangGraph的核心概念、任务规划策略和工具调用链构建展开实战。
LangGraph工作流编排核心概念
LangGraph将工作流抽象为有向图结构,每个节点是一个执行单元(函数或可调用对象),边定义了执行顺序和条件跳转逻辑。与线性链式调用不同,LangGraph支持循环和条件路由,这让Agent能够在执行过程中根据中间结果动态决定下一步操作。
核心组件包括三个部分:
State:贯穿整个工作流的共享状态对象,通常使用TypedDict定义,每个节点可以读取和更新其中的字段。
Node:工作流中的执行节点,接收当前状态作为输入,返回需要更新的状态字段。
Edge:连接节点的边,可以是固定路由(从A到B)或条件路由(根据状态内容选择下一个节点)。
from typing import TypedDict, Annotated, List
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
task_complete: bool
def call_model(state: AgentState) -> dict:
messages = state["messages"]
response = llm.invoke(messages)
return {"messages": [response]}
def should_continue(state: AgentState) -> str:
if state["task_complete"] or state["iteration"] >= 10:
return END
return "tools"
def call_tools(state: AgentState) -> dict:
last_message = state["messages"][-1]
tool_calls = last_message.tool_calls
results = []
for tc in tool_calls:
result = tool_map[tc["name"]].invoke(tc["args"])
results.append(result)
return {
"messages": results,
"tool_results": [str(r) for r in results],
"iteration": state["iteration"] + 1
}
多步骤任务规划与状态管理
实际业务场景中,Agent通常需要执行多个步骤才能完成复杂任务。合理的任务规划策略决定了工作流的效率和可靠性。一种常见做法是将任务拆分为”规划-执行-验证”三个阶段,规划阶段让模型分解任务目标,执行阶段调用工具逐步完成,验证阶段检查结果是否满足要求。
状态管理是工作流编排的基础。在LangGraph中,状态通过Reducer函数实现增量更新。使用operator.add作为Reducer时,列表类型字段会在每次节点执行后追加而非覆盖,这在维护消息历史时非常有用。
class PlanningState(TypedDict):
messages: Annotated[List[BaseMessage], operator.add]
plan: List[str]
current_step: int
results: Annotated[List[str], operator.add]
status: str
def planning_node(state: PlanningState) -> dict:
task_desc = state["messages"][-1].content
prompt = "根据以下任务制定执行计划:\n" + task_desc + "\n输出JSON格式的步骤列表"
response = llm.invoke([HumanMessage(content=prompt)])
plan = json.loads(response.content)
return {"plan": plan, "current_step": 0, "status": "executing"}
def execution_node(state: PlanningState) -> dict:
step_idx = state["current_step"]
step = state["plan"][step_idx]
result = execute_step(step)
return {
"results": [result],
"current_step": step_idx + 1
}
def verification_node(state: PlanningState) -> dict:
if state["current_step"] >= len(state["plan"]):
return {"status": "completed"}
return {"status": "executing"}
workflow = StateGraph(PlanningState)
workflow.add_node("planner", planning_node)
workflow.add_node("executor", execution_node)
workflow.add_node("verifier", verification_node)
workflow.set_entry_point("planner")
workflow.add_edge("planner", "executor")
workflow.add_edge("executor", "verifier")
workflow.add_conditional_edges(
"verifier",
lambda state: "executor" if state["status"] == "executing" else END
)
app = workflow.compile()
工具调用链构建实战
工具调用是Agent与外部系统交互的核心能力。构建工具调用链时,需要考虑工具注册、参数校验、错误处理和结果解析四个环节。LangGraph中工具节点通常与模型节点配合使用,模型决定调用哪些工具,工具节点执行实际操作并将结果回传给模型。
错误处理是工具调用链中容易被忽略的部分。网络超时、API限流、参数格式错误等问题需要在工作流层面统一处理,避免单个工具失败导致整个工作流中断。
from langchain_core.tools import tool
import httpx
import asyncio
@tool
def search_api(query: str, limit: int = 10) -> str:
"""搜索API工具,返回匹配结果"""
try:
resp = httpx.get(
"https://api.example.com/search",
params={"q": query, "limit": limit},
timeout=10
)
resp.raise_for_status()
return resp.text
except httpx.TimeoutException:
return json.dumps({"error": "请求超时", "retryable": True})
except httpx.HTTPStatusError as e:
if e.response.status_code == 429:
return json.dumps({"error": "API限流", "retryable": True})
return json.dumps({"error": str(e), "retryable": False})
@tool
def database_query(sql: str) -> str:
"""数据库查询工具"""
if not sql.strip().upper().startswith("SELECT"):
return json.dumps({"error": "仅支持SELECT查询"})
try:
rows = db.execute(sql).fetchall()
return json.dumps(rows, default=str)
except Exception as e:
return json.dumps({"error": str(e), "retryable": False})
def tool_node_with_retry(state: AgentState, max_retries: int = 2) -> dict:
last_message = state["messages"][-1]
results = []
for tc in last_message.tool_calls:
tool = tool_map.get(tc["name"])
if not tool:
results.append("未知工具: " + tc["name"])
continue
for attempt in range(max_retries + 1):
try:
result = tool.invoke(tc["args"])
results.append(str(result))
break
except Exception as e:
if attempt < max_retries:
import time
time.sleep(2 ** attempt)
else:
results.append("工具执行失败: " + str(e))
return {"messages": [HumanMessage(content="\n".join(results))]}
工作流持久化与人机协作
长时间运行的Agent工作流需要状态持久化支持。LangGraph通过checkpointer机制实现工作流状态的保存和恢复,使用MemorySaver进行内存存储,或使用SqliteSaver持久化到数据库。持久化后,工作流可以在任意节点暂停、恢复,也支持人机协作模式——在关键决策点暂停等待人工确认后再继续执行。
from langgraph.checkpoint.memory import MemorySaver
checkpointer = MemorySaver()
app = workflow.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "task-001"}}
# 执行工作流,可在任意节点中断
for event in app.stream(initial_state, config=config):
print(f"节点输出: {event}")
# 恢复执行
for event in app.stream(None, config=config):
print(f"恢复执行: {event}")
# 人机协作:在关键节点暂停
workflow.add_node("human_review", human_approval_node)
workflow.add_edge("executor", "human_review")
app = workflow.compile(
checkpointer=checkpointer,
interrupt_before=["human_review"]
)
LangGraph的循环执行和条件路由能力使其在构建ReAct风格Agent、多Agent协作系统、复杂数据处理流水线等场景中具备明显优势。实际部署中需关注工作流超时控制、并发执行管理和状态存储选型,根据业务复杂度选择合适的checkpoint存储方案。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/da-mo-xing-agent-gong-zuo-liu-bian-pai-shi-zhan-langgraph/