单个 Agent 处理复杂任务时上下文爆炸、工具冲突、无法并行。本文用完整可运行代码演示如何用 LangGraph 构建多Agent协作系统:Supervisor 调度 + Worker 分工 + 共享状态管理,从架构设计到生产部署全流程覆盖。
为什么需要多Agent编排
单个 LLM Agent 在处理复杂任务时面临三大瓶颈:
- 上下文窗口限制:一个 Agent 承载所有工具的描述和历史对话,很快耗尽 token 预算
- 工具选择混乱:注册 20+ 工具后,模型工具调用准确率显著下降
- 无法并行:串行执行无法利用子任务间的独立性
LangGraph 通过有向图建模 Agent 协作流程,支持条件分支、循环、人工介入和并行执行,是目前 Multi-Agent 工程化最成熟的框架之一。
LangGraph 核心概念
# LangGraph 的三个核心抽象:
# 1. State — 所有节点共享的不可变状态(TypedDict 或 Pydantic Model)
# 2. Node — 接收状态、执行逻辑、返回状态更新的函数
# 3. Edge — 连接节点的边,支持条件路由
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, END
from langgraph.graph.message import add_messages
# 定义共享状态
class AgentState(TypedDict):
messages: Annotated[list, add_messages] # 消息列表,自动累加
next_agent: str # 下一个要执行的 Agent
task_complete: bool # 任务是否完成
results: dict # 各 Agent 的结果汇总
架构设计:Supervisor + Worker 模式
我们构建一个「技术研究助手」系统,包含以下角色:
- Supervisor(调度者):分析用户需求,决定派发给哪个 Worker,汇总最终结果
- Researcher(研究员):使用搜索工具收集信息
- Coder(程序员):编写和执行代码
- Writer(撰稿人):将研究成果整理成报告
用户输入 → Supervisor → [Researcher | Coder | Writer] → Supervisor → 汇总输出
↑__________循环__________↓
完整实现
1. 定义状态与工具
import operator
from typing import TypedDict, Annotated, Literal
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage
from langchain_openai import ChatOpenAI
from langchain_core.tools import tool
from langgraph.graph import StateGraph, END, START
from langgraph.prebuilt import create_react_agent
# ============ 共享状态 ============
class TeamState(TypedDict):
messages: Annotated[list[BaseMessage], operator.add]
team_members: list[str]
next: str # Supervisor 决定的下一个 Agent
task: str # 原始任务描述
research_notes: str # 研究笔记
code_output: str # 代码执行结果
final_report: str # 最终报告
# ============ 工具定义 ============
@tool
def web_search(query: str) -> str:
"""搜索网络获取最新信息"""
# 实际项目中替换为 Tavily/SerperAPI
return f"搜索结果:关于 '{query}' 的最新信息显示..."
@tool
def execute_python(code: str) -> str:
"""执行 Python 代码并返回结果"""
try:
# 生产环境用沙箱执行
local_ns = {}
exec(code, {}, local_ns)
return str(local_ns.get("result", "执行成功,无返回值"))
except Exception as e:
return f"执行错误: {e}"
@tool
def write_report(topic: str, content: str) -> str:
"""将内容整理为结构化报告"""
return f"# 报告:{topic}\n\n{content}\n\n---\n报告生成完成。"
# ============ LLM 初始化 ============
llm = ChatOpenAI(model="gpt-4o", temperature=0)
2. 构建 Supervisor 节点
from langchain_core.prompts import ChatPromptTemplate
def supervisor_node(state: TeamState) -> dict:
"""Supervisor:分析当前状态,决定下一步派发给哪个Agent"""
system_prompt = """你是一个团队调度者。根据当前任务进度,决定下一步由哪个团队成员处理。
可选成员:
- researcher: 负责搜索和收集信息
- coder: 负责编写和执行代码
- writer: 负责整理报告
- FINISH: 任务已完成,输出最终结果
判断规则:
1. 如果还没有研究信息,派给 researcher
2. 如果需要代码验证,派给 coder
3. 如果研究和代码都完成,派给 writer
4. 如果报告已生成,返回 FINISH
"""
messages = [
{"role": "system", "content": system_prompt},
] + [msg for msg in state["messages"][-6:]] # 只看最近6条消息防止上下文爆炸
# 加入当前任务状态
context = f"任务: {state['task']}\n"
if state.get("research_notes"):
context += f"研究笔记: {state['research_notes'][:500]}\n"
if state.get("code_output"):
context += f"代码结果: {state['code_output'][:500]}\n"
if state.get("final_report"):
context += f"报告状态: 已生成\n"
messages.append({"role": "user", "content": context + "\n下一步派给谁?"})
response = llm.invoke(messages)
# 解析决策
decision = response.content.strip().lower()
if "finish" in decision:
next_agent = "FINISH"
elif "researcher" in decision:
next_agent = "researcher"
elif "coder" in decision:
next_agent = "coder"
elif "writer" in decision:
next_agent = "writer"
else:
next_agent = "FINISH" # 默认结束
return {
"next": next_agent,
"messages": [response],
}
3. 构建 Worker 节点
def create_worker_node(agent_name: str, tools: list, system_prompt: str):
"""创建一个 Worker Agent 节点"""
# 使用 LangGraph 的 ReAct Agent
agent = create_react_agent(llm, tools, prompt=system_prompt)
def worker_node(state: TeamState) -> dict:
# 构造给 Worker 的输入
task_context = f"原始任务: {state['task']}\n"
if state.get("research_notes") and agent_name != "researcher":
task_context += f"已有研究: {state['research_notes'][:1000]}\n"
if state.get("code_output") and agent_name != "coder":
task_context += f"已有代码结果: {state['code_output'][:1000]}\n"
result = agent.invoke({
"messages": [HumanMessage(content=task_context)]
})
# 提取最后一条 AI 消息
last_msg = result["messages"][-1]
# 更新对应字段
updates = {"messages": [last_msg]}
if agent_name == "researcher":
updates["research_notes"] = last_msg.content
elif agent_name == "coder":
updates["code_output"] = last_msg.content
elif agent_name == "writer":
updates["final_report"] = last_msg.content
return updates
return worker_node
# 创建三个 Worker
researcher_node = create_worker_node(
"researcher",
[web_search],
"你是研究员。使用 web_search 工具收集与任务相关的信息。"
"提供详细、准确的研究笔记。"
)
coder_node = create_worker_node(
"coder",
[execute_python],
"你是程序员。使用 execute_python 工具编写和执行代码。"
"代码中用 result 变量存储最终输出。"
)
writer_node = create_worker_node(
"writer",
[write_report],
"你是技术撰稿人。将研究和代码结果整理成结构清晰的 Markdown 报告。"
"使用 write_report 工具生成最终报告。"
)
4. 组装工作流图
def build_workflow():
"""构建 LangGraph 工作流"""
workflow = StateGraph(TeamState)
# 添加节点
workflow.add_node("supervisor", supervisor_node)
workflow.add_node("researcher", researcher_node)
workflow.add_node("coder", coder_node)
workflow.add_node("writer", writer_node)
# 设置入口
workflow.add_edge(START, "supervisor")
# Supervisor 的条件路由
def route_from_supervisor(state: TeamState) -> str:
next_agent = state["next"]
if next_agent == "FINISH":
return END
return next_agent
workflow.add_conditional_edges(
"supervisor",
route_from_supervisor,
{
"researcher": "researcher",
"coder": "coder",
"writer": "writer",
END: END,
},
)
# 所有 Worker 执行完后回到 Supervisor
workflow.add_edge("researcher", "supervisor")
workflow.add_edge("coder", "supervisor")
workflow.add_edge("writer", "supervisor")
# 编译
return workflow.compile()
# 构建并测试
app = build_workflow()
5. 运行工作流
# 初始状态
initial_state = {
"messages": [HumanMessage(content="分析 2026 年主流向量数据库的性能差异,并用代码对比 Milvus 和 Qdrant 的查询延迟")],
"team_members": ["researcher", "coder", "writer"],
"next": "",
"task": "分析 2026 年主流向量数据库的性能差异,并用代码对比 Milvus 和 Qdrant 的查询延迟",
"research_notes": "",
"code_output": "",
"final_report": "",
}
# 流式执行
for event in app.stream(initial_state, {"recursion_limit": 15}):
for node_name, output in event.items():
print(f"\n{'='*60}")
print(f"节点: {node_name}")
print(f"下一步: {output.get('next', 'N/A')}")
if "research_notes" in output and output["research_notes"]:
print(f"研究笔记: {output['research_notes'][:200]}...")
if "code_output" in output and output["code_output"]:
print(f"代码结果: {output['code_output'][:200]}...")
if "final_report" in output and output["final_report"]:
print(f"最终报告: {output['final_report'][:200]}...")
print("\n工作流执行完成")
高级特性
人工介入(Human-in-the-Loop)
from langgraph.checkpoint.memory import MemorySaver
# 添加检查点,支持暂停和人工审核
checkpointer = MemorySaver()
app_with_interrupt = build_workflow().compile(
checkpointer=checkpointer,
interrupt_before=["writer"], # 在写报告前暂停,等待人工确认
)
# 执行到 writer 前会暂停
config = {"configurable": {"thread_id": "thread-1"}}
for event in app_with_interrupt.stream(initial_state, config):
print(f"节点 {list(event.keys())} 完成")
# 人工审核后继续
user_feedback = input("是否继续生成报告?(y/n): ")
if user_feedback.lower() == "y":
for event in app_with_interrupt.stream(None, config):
print(f"节点 {list(event.keys())} 完成")
添加循环次数限制
def build_workflow_with_limit(max_iterations=10):
workflow = StateGraph(TeamState)
# 添加迭代计数器到状态
workflow.add_node("supervisor", supervisor_node)
workflow.add_node("researcher", researcher_node)
workflow.add_node("coder", coder_node)
workflow.add_node("writer", writer_node)
workflow.add_edge(START, "supervisor")
def route_from_supervisor(state: TeamState) -> str:
# 防止无限循环
msg_count = len(state["messages"])
if msg_count > max_iterations * 3: # 每轮约3条消息
print(f"⚠️ 达到最大迭代次数 {max_iterations},强制结束")
return END
next_agent = state["next"]
if next_agent == "FINISH":
return END
return next_agent
workflow.add_conditional_edges(
"supervisor", route_from_supervisor,
{"researcher": "researcher", "coder": "coder", "writer": "writer", END: END},
)
for worker in ["researcher", "coder", "writer"]:
workflow.add_edge(worker, "supervisor")
return workflow.compile()
生产部署
FastAPI 服务封装
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import json, asyncio
app_api = FastAPI(title="Multi-Agent API")
@app_api.post("/agent/run")
async def run_agent(task: str):
"""同步执行多Agent任务"""
state = {
"messages": [HumanMessage(content=task)],
"team_members": ["researcher", "coder", "writer"],
"next": "",
"task": task,
"research_notes": "",
"code_output": "",
"final_report": "",
}
result = app.invoke(state, {"recursion_limit": 15})
return {
"task": task,
"report": result.get("final_report", "未生成报告"),
"research": result.get("research_notes", ""),
"code": result.get("code_output", ""),
}
@app_api.post("/agent/stream")
async def stream_agent(task: str):
"""流式执行,实时返回各节点输出"""
async def event_stream():
state = {
"messages": [HumanMessage(content=task)],
"team_members": ["researcher", "coder", "writer"],
"next": "", "task": task,
"research_notes": "", "code_output": "", "final_report": "",
}
for event in app.stream(state, {"recursion_limit": 15}):
for node, output in event.items():
yield f"data: {json.dumps({'node': node, 'next': output.get('next', '')})}\n\n"
yield f"data: {json.dumps({'status': 'complete'})}\n\n"
return StreamingResponse(event_stream(), media_type="text/event-stream")
# 启动: uvicorn server:app_api --host 0.0.0.0 --port 8000
Docker 部署
FROM python:3.11-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
CMD ["uvicorn", "server:app_api", "--host", "0.0.0.0", "--port", "8000"]
# requirements.txt
langgraph>=0.2.0
langchain-openai>=0.2.0
langchain-core>=0.3.0
fastapi>=0.115.0
uvicorn>=0.32.0
常见问题 FAQ
Q1: Supervisor 总是选错 Agent 怎么办?
优化 Supervisor 的 system prompt,加入更明确的判断规则。也可以用结构化输出强制 LLM 返回 JSON 格式的决策:{"next": "researcher", "reason": "需要先收集信息"}。另外,确保传给 Supervisor 的状态摘要足够清晰。
Q2: 如何控制 Agent 之间的消息传递?
不要把所有消息都传给每个 Worker。在 Worker 节点中,只提取与该 Agent 相关的上下文(如任务描述 + 前序结果摘要)。用 state["messages"][-6:] 限制窗口大小,避免上下文爆炸。
Q3: 工作流卡在循环中出不来怎么办?
必须设置 recursion_limit(如 15),并在路由函数中加入迭代计数检查。如果 Supervisor 反复在两个 Worker 之间跳转,说明任务定义不够清晰或 Worker 输出不够明确,需要优化 prompt。
Q4: LangGraph 和 CrewAI 有什么区别?
LangGraph 基于状态图,控制流更精确,适合需要条件分支、循环和人工介入的复杂流程。CrewAI 基于角色和任务,更声明式,适合简单的线性协作。需要精细控制选 LangGraph,需要快速搭建选 CrewAI。
Q5: 如何接入 MCP 工具?
LangGraph 可通过 langchain-mcp-adapters 包接入 MCP Server。先创建 MCP Client 连接到 MCP Server,获取工具列表,然后绑定到 ReAct Agent。这样 Worker 可以使用任何 MCP 兼容的工具,实现工具生态复用。
Q6: 多 Agent 系统的成本怎么控制?
每个 Agent 调用都消耗 token。控制策略:1) Supervisor 用小模型(如 GPT-4o-mini),Worker 用大模型;2) 压缩传递给各 Agent 的上下文;3) 缓存工具调用结果避免重复搜索;4) 设置最大迭代次数。
总结
LangGraph 多Agent系统的核心是状态图 + 条件路由 + 共享状态。设计时遵循:Supervisor 做决策、Worker 做执行、状态做通信。生产部署注意三点:设循环上限防卡死、压缩上下文控成本、用检查点支持人工介入。当任务复杂到单个 Agent 无法有效处理时,多Agent编排是必然选择。