
当单一大模型无法满足复杂业务需求时,图结构编排的多智能体系统成为新的技术高地。本文将带你从原理到实践,完整实现一个可扩展的多智能体协作框架。
2025年已过半,LLM的能力边界正在被快速拓展。但一个残酷的事实是:单纯依赖Prompt工程和RAG,已经无法应对企业级场景的复杂需求。
我们面临的真实挑战包括:
LangGraph的出现,为这些问题提供了系统性的解决方案。
LangGraph的本质是一个基于图结构的状态编排引擎。它受Pregel计算模型启发,将Agent执行过程建模为:
State -> Node -> State -> Node -> State ...每个Node代表一个计算单元(可以是LLM调用、工具执行或自定义函数),Edge定义了状态流转的路径。
核心概念对应关系:
LangGraph概念 | 现实映射 |
|---|---|
StateGraph | 工作流蓝图 |
Node | 执行单元(Agent/工具/函数) |
Edge | 状态转移条件 |
State | 全局上下文对象 |
Checkpointer | 状态持久化与恢复 |
LangGraph的状态管理采用Reducer模式:
from typing import TypedDict, Annotated
from operator import add
class AgentState(TypedDict):
messages: Annotated[list, add] # 累加式更新
current_step: str # 覆盖式更新
tool_results: dict # 字典合并这种设计的精妙之处在于:
我们将构建一个具备以下能力的客服系统:
依赖版本:
langgraph>=0.2.0
langchain>=0.3.0
langchain-openai>=0.2.0状态定义:
from typing import TypedDict, Annotated, Literal
from operator import add
from langchain_core.messages import BaseMessage
class CustomerServiceState(TypedDict):
messages: Annotated[list[BaseMessage], add]
intent: str | None
entities: dict
retrieval_docs: list
need_human: bool
order_id: str | None节点实现:
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolNode
from langchain_openai import ChatOpenAI
from langchain_community.tools import tool
import json
# 初始化模型
model = ChatOpenAI(model="gpt-4", temperature=0)
# 1. 意图识别节点
def intent_classifier(state: CustomerServiceState) -> dict:
"""识别用户意图并提取实体"""
last_msg = state["messages"][-1].content
prompt = f"""
分析用户消息,返回JSON格式结果:
消息:{last_msg}
意图类型:order_query / product_consult / complaint / human_required
实体:{{"order_id": "", "product": ""}}
"""
response = model.invoke(prompt)
result = json.loads(response.content)
return {
"intent": result.get("intent"),
"entities": result.get("entities", {})
}
# 2. RAG检索节点
def retrieval_node(state: CustomerServiceState) -> dict:
"""从知识库检索相关信息"""
from langchain_community.vectorstores import Chroma
from langchain_openai import OpenAIEmbeddings
# 初始化向量库(实际使用时需持久化)
vectorstore = Chroma(
embedding_function=OpenAIEmbeddings(),
persist_directory="./kb_store"
)
query = state["messages"][-1].content
docs = vectorstore.similarity_search(query, k=3)
return {"retrieval_docs": docs}
# 3. 订单查询工具
@tool
def query_order(order_id: str) -> dict:
"""模拟订单查询API"""
# 实际场景中替换为真实API调用
mock_db = {
"ORD-001": {"status": "shipped", "date": "2025-08-20"},
"ORD-002": {"status": "pending", "date": "2025-08-22"}
}
return mock_db.get(order_id, {"error": "订单不存在"})
# 4. 响应生成节点
def response_generator(state: CustomerServiceState) -> dict:
"""基于当前状态生成最终回复"""
context = ""
if state.get("retrieval_docs"):
context = "\n".join([d.page_content for d in state["retrieval_docs"]])
if state.get("intent") == "order_query" and state.get("entities", {}).get("order_id"):
order_result = query_order(state["entities"]["order_id"])
context += f"\n订单信息:{json.dumps(order_result)}"
prompt = f"""
你是一个专业的客服助手,基于以下信息回复用户:
上下文:{context}
用户最新消息:{state['messages'][-1].content}
注意:
- 如果用户要求转人工,引导用户提供联系方式
- 对于订单查询,直接告知状态信息
- 保持礼貌专业的语气
"""
response = model.invoke(prompt)
return {"messages": [response]}
# 5. 路由决策函数
def route_after_intent(state: CustomerServiceState) -> Literal["retrieval", "order_tool", "human_handoff", END]:
"""根据意图决定下一个节点"""
intent = state.get("intent")
if intent == "human_required":
return "human_handoff"
elif intent == "order_query":
return "order_tool"
elif intent in ["product_consult", "complaint"]:
return "retrieval"
else:
return END图编排:
# 构建状态图
workflow = StateGraph(CustomerServiceState)
# 添加节点
workflow.add_node("classify", intent_classifier)
workflow.add_node("retrieval", retrieval_node)
workflow.add_node("order_tool", lambda state: {"messages": [query_order(state["entities"].get("order_id", ""))]})
workflow.add_node("human_handoff", lambda state: {"need_human": True, "messages": ["正在为您转接人工..."]})
workflow.add_node("generate", response_generator)
# 设置入口
workflow.set_entry_point("classify")
# 添加条件边
workflow.add_conditional_edges(
"classify",
route_after_intent,
{
"retrieval": "retrieval",
"order_tool": "order_tool",
"human_handoff": "human_handoff",
END: END
}
)
# 后续边连接
workflow.add_edge("retrieval", "generate")
workflow.add_edge("order_tool", "generate")
workflow.add_edge("human_handoff", END)
workflow.add_edge("generate", END)
# 编译
app = workflow.compile()from langchain_core.messages import HumanMessage
# 执行示例
result = app.invoke({
"messages": [HumanMessage(content="我的订单ORD-001什么时候发货?")],
"intent": None,
"entities": {},
"retrieval_docs": [],
"need_human": False
})
print(result["messages"][-1].content)当单个Agent无法满足需求时,我们引入监督者-工作者模式:
from langgraph.graph import StateGraph
from langgraph.prebuilt import create_react_agent
# 创建专业Agent
researcher_agent = create_react_agent(model, tools=[search_tool, web_scraper])
coder_agent = create_react_agent(model, tools=[python_repl, git_tool])
analyst_agent = create_react_agent(model, tools=[data_visualizer, sql_tool])
class MultiAgentState(TypedDict):
messages: Annotated[list, add]
next_agent: str
task_completed: bool
def supervisor(state: MultiAgentState) -> dict:
"""监督者决策下一个执行者"""
prompt = f"""
当前任务进度:{state['messages']}
可选Agent:researcher, coder, analyst
请选择下一个应该执行的Agent,或返回"FINISH"。
"""
response = model.invoke(prompt)
return {"next_agent": response.content.strip()}
# 构建多Agent图
multi_workflow = StateGraph(MultiAgentState)
multi_workflow.add_node("supervisor", supervisor)
multi_workflow.add_node("researcher", researcher_agent)
multi_workflow.add_node("coder", coder_agent)
multi_workflow.add_node("analyst", analyst_agent)
# 循环执行直到完成
multi_workflow.add_conditional_edges(
"supervisor",
lambda s: s["next_agent"],
{
"researcher": "researcher",
"coder": "coder",
"analyst": "analyst",
"FINISH": END
}
)多Agent协作的关键是状态共享机制:
class SharedState(TypedDict):
task_description: str
research_findings: list
code_snippets: list
analysis_results: dict
current_phase: str
errors: list通过将状态设计为共享内存,每个Agent可以:
from langgraph.checkpoint import MemorySaver
import logging
# 启用状态持久化
memory = MemorySaver()
app = workflow.compile(checkpointer=memory)
# 添加回调监控
from langchain_core.callbacks import BaseCallbackHandler
class LoggingCallback(BaseCallbackHandler):
def on_chain_start(self, serialized, inputs, **kwargs):
logging.info(f"开始执行: {serialized.get('name')}")
def on_chain_end(self, outputs, **kwargs):
logging.info(f"执行完成: {outputs}")
# 执行时传入回调
result = app.invoke(
{"messages": [HumanMessage(content="...")]},
config={"callbacks": [LoggingCallback()]}
)from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10))
def robust_llm_call(prompt: str):
try:
return model.invoke(prompt)
except Exception as e:
logging.error(f"LLM调用失败: {e}")
raise# 流式执行
async for event in app.astream_events(
{"messages": [HumanMessage(content="...")]},
version="v2"
):
if event["event"] == "on_chat_model_stream":
print(event["data"]["chunk"].content, end="")add_node的branch参数实现并行执行LangGraph提供了一套完整的声明式Agent编排框架,其核心价值在于:
未来演进方向:
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。