首页 / AI工具 / AI工作流自动化实战:LangGraph 多Agent协作编排从设计到部署...

AI工作流自动化实战:LangGraph 多Agent协作编排从设计到部署

单个 Agent 处理复杂任务时上下文爆炸、工具冲突、无法并行。本文用完整可运行代码演示如何用 LangGraph 构建多Agent协作系统:Supervisor 调度 + Worker 分工 + 共享状态管理,从架构设计到生产部署全流程覆盖。

为什么需要多Agent编排

单个 LLM Agent 在处理复杂任务时面临三大瓶颈:

  1. 上下文窗口限制:一个 Agent 承载所有工具的描述和历史对话,很快耗尽 token 预算
  2. 工具选择混乱:注册 20+ 工具后,模型工具调用准确率显著下降
  3. 无法并行:串行执行无法利用子任务间的独立性

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 → [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编排是必然选择。