AI Agent工作流编排实战:LangGraph状态图与Function Calling工具链路构建

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/

(0)
小编小编
上一篇 9小时前
下一篇 8小时前

相关推荐