作者:PySuper | 来源:zhengxingtao.com | 更新日期:2026-11-15
关联阅读:第19篇 LangGraph StateGraph 状态机设计 | 第21篇 多Agent协作编排理论
写在前面
上一篇我们聊了多 Agent 协作编排的理论框架,这篇终于到了工程落地环节。如果你已经用 StateGraph 写过几十个节点的单体图,你一定体会过那种"改一个节点,牵一发动全身"的恐惧——这就是我们今天要解决的问题。
LangGraph 的子图(Subgraph)机制和多 Agent 编排模式,本质上是在回答一个工程问题:怎么让 LLM 应用在变复杂的同时,还能被人类理解和维护?
本文基于 langgraph 0.2.x,所有代码可运行。Let's go.
一、为什么需要子图
单图膨胀:50+ 节点的维护噩梦
当你用 StateGraph 构建复杂应用时,图会以一种让人绝望的速度膨胀:
plaintext
┌─────────────────────────────────────────────────────────────┐
│ 一个真实的研究助手图 │
│ │
│ START → 接收问题 → 意图分类 → 搜索A → 搜索B → 搜索C │
│ → 摘要 → 评估 → [不够?] → 搜索A ... │
│ → 写作 → 审核 → [不合格?] → 修改 → 写作 ... │
│ → 格式化 → 输出 → END │
│ │
│ 50+ 节点,80+ 条边,一个人维护已经疯了 │
└─────────────────────────────────────────────────────────────┘
这不是夸张,我见过生产环境里一个图 60 多个节点,光是理解节点间的数据流向就要半天。
子图的核心价值
子图把大问题拆成小问题,每个小问题独立演进:
表格
一句话总结:子图让你像搭乐高一样组装 AI 工作流,而不是在面团上越揉越大。
二、Subgraph 基础
子图定义方式
LangGraph 中有两种方式让子图与父图通信:
plaintext
方式1: 共享状态键(Shared Keys) 方式2: 状态转换(Transform)
───────────────────────── ──────────────────────────
父图 State: {foo, bar} 父图 State: {foo}
子图 State: {foo, bar, baz} 子图 State: {qux, quux}
↕ 共享 foo, bar ↕ 手动转换
子图直接作为 add_node 参数 包装函数做 foo↔qux 映射
方式 1:共享状态键 —— 最简单,子图直接作为节点
当子图和父图有共同的状态键时,直接把编译后的子图传给 add_node:
python
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
# 父子图共享同一个 State 类型
class State(TypedDict):
foo: str
# ---- 子图 ----
def subgraph_node_1(state: State):
return {"foo": "hi! " + state["foo"]}
subgraph_builder = StateGraph(State)
subgraph_builder.add_node(subgraph_node_1)
subgraph_builder.add_edge(START, "subgraph_node_1")
subgraph = subgraph_builder.compile()
# ---- 父图 ----
builder = StateGraph(State)
builder.add_node("node_1", subgraph) # ← 编译后的子图直接作为节点
builder.add_edge(START, "node_1")
graph = builder.compile()
# 运行
result = graph.invoke({"foo": "world"})
print(result["foo"]) # "hi! world"
方式 2:状态转换 —— 更灵活,子图有独立状态
当子图和父图的状态键完全不同时,你需要写一个包装函数做转换:
python
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START
# ---- 子图状态(与父图没有共同键) ----
class SubgraphState(TypedDict):
bar: str # 私有键
baz: str # 私有键
def subgraph_node_1(state: SubgraphState):
return {"baz": "baz"}
def subgraph_node_2(state: SubgraphState):
return {"bar": state["bar"] + state["baz"]}
subgraph_builder = StateGraph(SubgraphState)
subgraph_builder.add_node(subgraph_node_1)
subgraph_builder.add_node(subgraph_node_2)
subgraph_builder.add_edge(START, "subgraph_node_1")
subgraph_builder.add_edge("subgraph_node_1", "subgraph_node_2")
subgraph = subgraph_builder.compile()
# ---- 父图 ----
class ParentState(TypedDict):
foo: str
def node_1(state: ParentState):
return {"foo": "hi! " + state["foo"]}
def node_2(state: ParentState):
# 入场:父图状态 → 子图状态
response = subgraph.invoke({"bar": state["foo"]})
# 出场:子图状态 → 父图状态
return {"foo": response["bar"]}
builder = StateGraph(ParentState)
builder.add_node("node_1", node_1)
builder.add_node("node_2", node_2) # ← 用包装函数代替子图
builder.add_edge(START, "node_1")
builder.add_edge("node_1", "node_2")
graph = builder.compile()
result = graph.invoke({"foo": "foo"})
print(result["foo"]) # "hi! foobaz"
父图与子图的状态传递机制
用一张图说清楚两种模式:
plaintext
┌─────────────────────── 共享键模式 ───────────────────────┐
│ │
│ ParentState: {foo: str, bar: str} │
│ │ │
│ │ foo, bar 自动传递(共享键) │
│ ▼ │
│ SubgraphState: {foo: str, bar: str, baz: str} │
│ │ ↑ 私有键,父图不可见 │
│ │ foo, bar 自动回传 │
│ ▼ │
│ ParentState: {foo: "updated", bar: "updated"} │
│ │
└──────────────────────────────────────────────────────────┘
┌─────────────────────── 转换模式 ─────────────────────────┐
│ │
│ ParentState: {foo: str} │
│ │ │
│ │ 手动映射: {"bar": state["foo"]} │
│ ▼ │
│ SubgraphState: {bar: str, baz: str} │
│ │ │
│ │ 手动映射: {"foo": response["bar"]} │
│ ▼ │
│ ParentState: {foo: "result from subgraph"} │
│ │
└──────────────────────────────────────────────────────────┘
状态键映射规则
共享键模式下,核心规则只有三条:
子图的 State 必须包含所有与父图共享的键,子图可以有额外私有键
共享键上的 Reducer 必须一致(父子图的同一个键用同一个 Reducer)
子图返回时,只有共享键的值会传回父图,私有键留在子图内部
三、子图的状态隔离
私有状态 vs 共享状态
这是子图最核心的设计:父图只能看到子图的输入/输出,看不到内部状态。
plaintext
┌─ 父图视角 ──────────────────────────────────────────────┐
│ │
│ 子图节点 │
│ ┌─────────────────────────────────────────┐ │
│ │ 输入: {foo} 输出: {foo, bar} │ │
│ │ ┌─────────────────────────────────┐ │ │
│ │ │ 内部状态 baz = "..." │ │ │
│ │ │ 内部状态 temp = [...] │ │ ← 父图看不见 │
│ │ │ 中间计算结果 │ │ │
│ │ └─────────────────────────────────┘ │ │
│ └─────────────────────────────────────────┘ │
│ │
│ 父图只关心:你给我什么,你给我什么 │
│ │
└─────────────────────────────────────────────────────────┘
带私有状态的子图:实战示例
来看一个日志分析系统,父图是入口,两个子图分别做"问题摘要"和"失败分析":
python
from typing import Optional, Annotated
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
# ---- 日志结构 ----
class Logs(TypedDict):
id: str
question: str
answer: str
grade: Optional[int]
feedback: Optional[str]
# 自定义 Reducer:按 id 去重合并
def add_logs(left: list[Logs], right: list[Logs]) -> list[Logs]:
if not left:
left = []
if not right:
right = []
logs = left.copy()
left_id_to_idx = {log["id"]: idx for idx, log in enumerate(logs)}
for log in right:
idx = left_id_to_idx.get(log["id"])
if idx is not None:
logs[idx] = log # 已存在则更新
else:
logs.append(log) # 不存在则追加
return logs
# ---- 失败分析子图 ----
class FailureAnalysisState(TypedDict):
logs: Annotated[list[Logs], add_logs] # 共享键
failure_report: str # 共享键(输出给父图)
failures: list[Logs] # ← 私有键!父图看不到
def get_failures(state: FailureAnalysisState):
failures = [log for log in state["logs"] if log.get("grade") == 0]
return {"failures": failures} # 写入私有键
def generate_failure_summary(state: FailureAnalysisState):
failure_ids = [log["id"] for log in state["failures"]]
return {"failure_report": f"Poor quality for doc IDs: {', '.join(failure_ids)}"}
fa_builder = StateGraph(FailureAnalysisState)
fa_builder.add_node("get_failures", get_failures)
fa_builder.add_node("generate_summary", generate_failure_summary)
fa_builder.add_edge(START, "get_failures")
fa_builder.add_edge("get_failures", "generate_summary")
fa_builder.add_edge("generate_summary", END)
fa_subgraph = fa_builder.compile()
# ---- 问题摘要子图 ----
class QuestionSummarizationState(TypedDict):
logs: Annotated[list[Logs], add_logs] # 共享键
summary_report: str # 共享键(输出给父图)
summary: str # ← 私有键!
def generate_question_summary(state: QuestionSummarizationState):
return {"summary": "Questions focused on ChatOllama and Chroma."}
def send_to_slack(state: QuestionSummarizationState):
return {"summary_report": state["summary"]}
qs_builder = StateGraph(QuestionSummarizationState)
qs_builder.add_node("generate_summary", generate_question_summary)
qs_builder.add_node("send_to_slack", send_to_slack)
qs_builder.add_edge(START, "generate_summary")
qs_builder.add_edge("generate_summary", "send_to_slack")
qs_builder.add_edge("send_to_slack", END)
qs_subgraph = qs_builder.compile()
# ---- 父图(入口图) ----
class EntryGraphState(TypedDict):
raw_logs: Annotated[list[Logs], add_logs]
logs: Annotated[list[Logs], add_logs] # 传给子图
failure_report: str # 从子图回收
summary_report: str # 从子图回收
def select_logs(state: EntryGraphState):
return {"logs": [log for log in state["raw_logs"] if "grade" in log]}
entry_builder = StateGraph(EntryGraphState)
entry_builder.add_node("select_logs", select_logs)
entry_builder.add_node("failure_analysis", fa_subgraph) # 子图作为节点
entry_builder.add_node("question_summarization", qs_subgraph) # 子图作为节点
entry_builder.add_edge(START, "select_logs")
entry_builder.add_edge("select_logs", "failure_analysis")
entry_builder.add_edge("select_logs", "question_summarization")
entry_builder.add_edge("failure_analysis", END)
entry_builder.add_edge("question_summarization", END)
graph = entry_builder.compile()
注意:failures 和 summary 是子图私有状态,父图完全不知道它们的存在。这就是状态隔离的力量。
⚠️ 大坑:operator.add 的双重拼接问题
这是子图里最容易踩的坑。当共享键使用了 operator.add 作为 Reducer,且子图作为节点直接 add_node 时,子图的输入状态已经包含了父图传过来的值,子图执行完后返回的值又会和父图的值再拼一次,导致双重拼接。
plaintext
┌─ 问题演示 ──────────────────────────────────────────────┐
│ │
│ 父图 path = ["grandparent", "parent"] │
│ │ │
│ │ 传入子图 │
│ ▼ │
│ 子图读取 path = ["grandparent", "parent"] │
│ 子图执行: path += ["child_start", "child_end"] │
│ 子图返回: ["grandparent", "parent", │
│ "child_start", "child_end"] │
│ │ │
│ │ Reducer (operator.add) 再拼一次! │
│ ▼ │
│ 父图 path = ["grandparent", "parent", │
│ "grandparent", "parent", │
│ "child_start", "child_end"] │
│ ↑ 重复了! │
│ │
└─────────────────────────────────────────────────────────┘
解法 1:自定义 Reducer 去重
python
import uuid
from typing import Annotated
from typing_extensions import TypedDict
def reduce_unique(left: list | None, right: list | None) -> list:
"""带 id 去重的 reducer,避免双重拼接"""
if not left:
left = []
if not right:
right = []
left_, right_ = [], []
for orig, new in [(left, left_), (right, right_)]:
for val in orig:
if not isinstance(val, dict):
val = {"val": val}
if "id" not in val:
val["id"] = str(uuid.uuid4())
new.append(val)
# 按 id 合并,有则替换,无则追加
left_idx_by_id = {val["id"]: i for i, val in enumerate(left_)}
merged = left_.copy()
for val in right_:
idx = left_idx_by_id.get(val["id"])
if idx is not None:
merged[idx] = val
else:
merged.append(val)
return merged
class ChildState(TypedDict):
name: str
path: Annotated[list[str], reduce_unique] # ← 用自定义 Reducer
class ParentState(TypedDict):
name: str
path: Annotated[list[str], reduce_unique] # ← 父子图用同一个 Reducer
解法 2:用包装函数代替直接 add_node
python
def call_child_graph(state: ParentState):
# 只传子图需要的键,控制输出
child_output = child_graph.invoke({"name": state["name"]})
return {"name": child_output["name"]} # 只取需要的字段
解法 3:子图状态键和父图完全分离,用转换模式
这样就没有共享键,自然没有双重拼接问题。
实战建议:如果子图和父图都有
Annotated[list, operator.add]的共享键,优先用解法 1 或解法 2,别指望裸operator.add能正常工作。
四、实战:可复用的研究子图
来一个完整的、可运行的研究子图。这个子图的设计思路是:搜索 → 摘要 → 评估 → 循环/退出,它内部有自己的循环逻辑,但从主图看只是一个"黑盒"节点。
python
import operator
from typing import Annotated, Literal
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
# ---- 研究子图状态 ----
class ResearchState(TypedDict):
query: str
search_results: Annotated[list[str], operator.add]
summary: str
evaluation: str
iteration: int
max_iterations: int
# 模拟搜索
def search_node(state: ResearchState):
# 实际项目里这里接 Tavily / SerpAPI 等
result = f"[Search round {state['iteration']}] Found results for: {state['query']}"
return {"search_results": [result]}
# 模拟摘要
def summarize_node(state: ResearchState):
results_text = "\n".join(state["search_results"])
summary = f"Summary of {len(state['search_results'])} results for '{state['query']}'"
return {"summary": summary}
# 模拟评估
def evaluate_node(state: ResearchState):
if state["iteration"] >= state["max_iterations"]:
return {"evaluation": "sufficient"}
return {"evaluation": "need_more"}
# 路由:评估通过则退出,否则继续搜索
def should_continue(state: ResearchState) -> Literal["search", "exit"]:
if state["evaluation"] == "sufficient":
return "exit"
return "search"
# 更新迭代计数
def increment_iteration(state: ResearchState):
return {"iteration": state["iteration"] + 1}
# ---- 构建研究子图 ----
research_builder = StateGraph(ResearchState)
research_builder.add_node("search", search_node)
research_builder.add_node("summarize", summarize_node)
research_builder.add_node("evaluate", evaluate_node)
research_builder.add_node("increment", increment_iteration)
research_builder.add_edge(START, "search")
research_builder.add_edge("search", "summarize")
research_builder.add_edge("summarize", "evaluate")
research_builder.add_conditional_edges("evaluate", should_continue, {
"search": "increment",
"exit": END,
})
research_builder.add_edge("increment", "search")
research_subgraph = research_builder.compile()
# ---- 主图:研究 → 写作 → 审核 ----
class MainState(TypedDict):
topic: str
research_summary: str
article: str
review: str
def call_research(state: MainState):
# 调用研究子图,做状态转换
result = research_subgraph.invoke({
"query": state["topic"],
"search_results": [],
"summary": "",
"evaluation": "",
"iteration": 0,
"max_iterations": 2,
})
return {"research_summary": result["summary"]}
def write_node(state: MainState):
article = f"Based on research: {state['research_summary']}, here is the article."
return {"article": article}
def review_node(state: MainState):
review = "APPROVED" if len(state["article"]) > 20 else "NEEDS_REVISION"
return {"review": review}
# ---- 构建主图 ----
builder = StateGraph(MainState)
builder.add_node("research", call_research) # 子图通过包装函数接入
builder.add_node("write", write_node)
builder.add_node("review", review_node)
builder.add_edge(START, "research")
builder.add_edge("research", "write")
builder.add_edge("write", "review")
builder.add_edge("review", END)
main_graph = builder.compile()
# ---- 运行 ----
result = main_graph.invoke({"topic": "LangGraph subgraphs", "research_summary": "", "article": "", "review": ""})
print(result["review"]) # APPROVED
print(result["article"][:80]) # Based on research: Summary of 2 results for 'LangGraph subgraphs'...
看,研究子图内部有循环(搜索→摘要→评估→可能再来一轮),但主图只看到:输入 topic,输出 research_summary。这就是子图的封装威力。
五、多Agent编排模式
5.1 Supervisor 模式(最常用)
Supervisor 模式的核心思想:一个总控 Agent 路由任务到专业 Agent,每个 Agent 是一个子图。
plaintext
┌──────────────── Supervisor 模式 ────────────────┐
│ │
│ 用户请求 │
│ │ │
│ ▼ │
│ ┌──────────────┐ │
│ │ Supervisor │ ← 只做路由 + 协调 + 整合 │
│ └──────┬───────┘ │
│ │ tool_call │
│ ┌────┼────┐ │
│ ▼ ▼ ▼ │
│ [研究] [写作] [审核] ← 专业 Agent(子图) │
│ │ │ │ │
│ └────┼────┘ │
│ ▼ │
│ Supervisor 整合 → 最终答案 │
│ │
└──────────────────────────────────────────────────┘
完整代码:研究 + 写作 + 审核 三 Agent 协作
python
import operator
from typing import Annotated, Literal
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langchain_core.messages import HumanMessage, SystemMessage, AIMessage
from langchain_openai import ChatOpenAI
# ---- 全局状态 ----
class AgentState(TypedDict):
messages: Annotated[list, add_messages]
next_agent: str
research_result: str
article: str
review_result: str
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0)
# ---- 研究 Agent(子图) ----
class ResearchAgentState(TypedDict):
messages: Annotated[list, add_messages]
query: str
research_result: str
def research_think(state: ResearchAgentState):
"""模拟研究 Agent:分析问题并搜索"""
result = f"Research findings about: {state['query']}"
return {
"messages": [AIMessage(content=f"I found: {result}")],
"research_result": result,
}
research_builder = StateGraph(ResearchAgentState)
research_builder.add_node("think", research_think)
research_builder.add_edge(START, "think")
research_builder.add_edge("think", END)
research_agent = research_builder.compile()
# ---- 写作 Agent(子图) ----
class WritingAgentState(TypedDict):
messages: Annotated[list, add_messages]
research_result: str
article: str
def writing_think(state: WritingAgentState):
"""模拟写作 Agent:基于研究结果写文章"""
article = f"Article based on: {state['research_result']}"
return {
"messages": [AIMessage(content=f"I wrote: {article[:50]}...")],
"article": article,
}
writing_builder = StateGraph(WritingAgentState)
writing_builder.add_node("think", writing_think)
writing_builder.add_edge(START, "think")
writing_builder.add_edge("think", END)
writing_agent = writing_builder.compile()
# ---- 审核 Agent(子图) ----
class ReviewAgentState(TypedDict):
messages: Annotated[list, add_messages]
article: str
review_result: str
def review_think(state: ReviewAgentState):
"""模拟审核 Agent:检查文章质量"""
verdict = "APPROVED" if len(state["article"]) > 20 else "NEEDS_REVISION"
return {
"messages": [AIMessage(content=f"Review: {verdict}")],
"review_result": verdict,
}
review_builder = StateGraph(ReviewAgentState)
review_builder.add_node("think", review_think)
review_builder.add_edge(START, "think")
review_builder.add_edge("think", END)
review_agent = review_builder.compile()
# ---- Supervisor 节点 ----
def supervisor(state: AgentState):
"""Supervisor:决定下一个该谁干活"""
# 简单路由逻辑(生产中用 LLM 做 tool_call 路由)
if not state.get("research_result"):
return {"next_agent": "research"}
elif not state.get("article"):
return {"next_agent": "writing"}
elif not state.get("review_result"):
return {"next_agent": "review"}
else:
return {"next_agent": "FINISH"}
# ---- 各 Agent 的包装函数 ----
def call_research(state: AgentState):
last_msg = state["messages"][-1].content if state["messages"] else ""
result = research_agent.invoke({"messages": [], "query": last_msg, "research_result": ""})
return {"research_result": result["research_result"]}
def call_writing(state: AgentState):
result = writing_agent.invoke({
"messages": [],
"research_result": state["research_result"],
"article": "",
})
return {"article": result["article"]}
def call_review(state: AgentState):
result = review_agent.invoke({
"messages": [],
"article": state["article"],
"review_result": "",
})
return {"review_result": result["review_result"]}
# 路由函数
def route_agent(state: AgentState) -> Literal["research", "writing", "review", "__end__"]:
if state["next_agent"] == "FINISH":
return "__end__"
return state["next_agent"]
# ---- 构建主图 ----
builder = StateGraph(AgentState)
builder.add_node("supervisor", supervisor)
builder.add_node("research", call_research)
builder.add_node("writing", call_writing)
builder.add_node("review", call_review)
builder.add_edge(START, "supervisor")
builder.add_conditional_edges("supervisor", route_agent, {
"research": "research",
"writing": "writing",
"review": "review",
"__end__": END,
})
builder.add_edge("research", "supervisor") # 回到 Supervisor
builder.add_edge("writing", "supervisor")
builder.add_edge("review", "supervisor")
supervisor_graph = builder.compile()
# ---- 运行 ----
result = supervisor_graph.invoke({
"messages": [HumanMessage(content="Write about LangGraph subgraphs")],
"next_agent": "",
"research_result": "",
"article": "",
"review_result": "",
})
print(result["review_result"]) # APPROVED
5.2 并行召回模式
多个检索通道并行执行,用 Reducer 合并去重。这比串行检索快 3-5 倍:
plaintext
┌───────────────── 并行召回 ─────────────────┐
│ │
│ 用户查询 │
│ │ │
│ ┌─────┼─────┐ │
│ ▼ ▼ ▼ │
│ [向量库] [关键词] [知识图谱] │
│ │ │ │ │
│ └─────┼─────┘ │
│ ▼ │
│ Reducer 合并去重 │
│ │ │
│ ▼ │
│ Rerank → 输出 │
│ │
└─────────────────────────────────────────────┘
python
import operator
from typing import Annotated
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
class RetrievalState(TypedDict):
query: str
# operator.add 让多个并行节点的结果自动拼接
docs: Annotated[list[str], operator.add]
reranked: list[str]
def vector_search(state: RetrievalState):
# 模拟向量库检索
return {"docs": [f"vec_doc_{state['query']}_1", f"vec_doc_{state['query']}_2"]}
def keyword_search(state: RetrievalState):
# 模拟关键词检索
return {"docs": [f"kw_doc_{state['query']}_1", f"vec_doc_{state['query']}_1"]} # 有重复!
def kg_search(state: RetrievalState):
# 模拟知识图谱检索
return {"docs": [f"kg_doc_{state['query']}_1"]}
def rerank(state: RetrievalState):
# 去重 + 排序
unique = list(dict.fromkeys(state["docs"])) # 保序去重
return {"reranked": unique}
# 构建并行召回图
builder = StateGraph(RetrievalState)
builder.add_node("vector_search", vector_search)
builder.add_node("keyword_search", keyword_search)
builder.add_node("kg_search", kg_search)
builder.add_node("rerank", rerank)
# 三个检索节点并行执行
builder.add_edge(START, "vector_search")
builder.add_edge(START, "keyword_search")
builder.add_edge(START, "kg_search")
# 汇聚到 rerank
builder.add_edge("vector_search", "rerank")
builder.add_edge("keyword_search", "rerank")
builder.add_edge("kg_search", "rerank")
builder.add_edge("rerank", END)
graph = builder.compile()
result = graph.invoke({"query": "LangGraph", "docs": [], "reranked": []})
print(result["reranked"])
# ['vec_doc_LangGraph_1', 'vec_doc_LangGraph_2', 'kw_doc_LangGraph_1', 'kg_doc_LangGraph_1']
# 重复的 vec_doc_LangGraph_1 被去重了
5.3 Map-Reduce 模式
当子任务数量在运行时才确定,用 Send API 实现动态路由:
plaintext
┌─────────────── Map-Reduce ───────────────┐
│ │
│ 生成主题列表(数量未知) │
│ │ │
│ ▼ │
│ ┌─ Send("joke", {subject: "lions"}) ┐ │
│ │ Send("joke", {subject: "cats"}) │ │
│ │ Send("joke", {subject: "Python"}) │ │ ← 动态并行
│ └──────────────────────────────────────┘ │
│ │ │
│ ▼ │
│ 每个 Send 实例独立执行 │
│ 结果通过 Reducer (operator.add) 自动合并 │
│ │ │
│ ▼ │
│ 选择最佳 → END │
│ │
└───────────────────────────────────────────┘
python
import operator
from typing import Annotated
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
# ---- 状态定义 ----
class OverallState(TypedDict):
topic: str
subjects: list[str]
jokes: Annotated[list[str], operator.add] # Reducer 自动合并
best_joke: str
class JokeState(TypedDict):
subject: str # 每个 Map 实例只接收一个 subject
# ---- 节点 ----
def generate_topics(state: OverallState):
"""生成主题列表(数量不确定)"""
topics = {
"animals": ["lions", "elephants", "penguins"],
"programming": ["Python", "JavaScript", "Rust"],
}
return {"subjects": topics.get(state["topic"], ["default"])}
def generate_joke(state: JokeState):
"""为单个主题生成笑话(每个 Send 实例独立执行)"""
joke_map = {
"lions": "Why don't lions like fast food? They can't catch it!",
"elephants": "Why don't elephants use computers? Afraid of the mouse!",
"penguins": "Why don't penguins talk to strangers? Hard to break the ice.",
"Python": "Why do Python devs prefer dark mode? Light attracts bugs!",
"JavaScript": "Why is JS like a car? It has NaN problems.",
"Rust": "Why do Rustaceans never get lost? The borrow checker!",
}
joke = joke_map.get(state["subject"], f"No joke for {state['subject']}")
return {"jokes": [joke]} # 注意:返回列表,Reducer 会拼接
def continue_to_jokes(state: OverallState):
"""路由函数:为每个 subject 创建一个 Send 实例"""
return [Send("generate_joke", {"subject": s}) for s in state["subjects"]]
def select_best(state: OverallState):
"""Reduce:选择最好的笑话"""
return {"best_joke": state["jokes"][0]} # 简化,实际用 LLM 选择
# ---- 构建图 ----
builder = StateGraph(OverallState)
builder.add_node("generate_topics", generate_topics)
builder.add_node("generate_joke", generate_joke)
builder.add_node("select_best", select_best)
builder.add_edge(START, "generate_topics")
# 条件边返回 Send 列表 → 动态并行
builder.add_conditional_edges("generate_topics", continue_to_jokes, ["generate_joke"])
builder.add_edge("generate_joke", "select_best")
builder.add_edge("select_best", END)
graph = builder.compile()
result = graph.invoke({"topic": "animals", "subjects": [], "jokes": [], "best_joke": ""})
print(result["jokes"])
# ['Why don't lions like fast food? They can't catch it!',
# "Why don't elephants use computers? Afraid of the mouse!",
# "Why don't penguins talk to strangers? Hard to break the ice."]
print(result["best_joke"])
# Why don't lions like fast food? They can't catch it!
Send API 的关键点:
表格
六、Deep Agents:规划 + 子Agent + 文件系统
LangGraph 0.2.x 引入了 Deep Agents 概念,这是一种更高级的多 Agent 架构:
plaintext
┌──────────────── Deep Agent 架构 ────────────────┐
│ │
│ ┌──────────────────────────────────────────┐ │
│ │ 主 Agent(Orchestrator) │ │
│ │ ┌────────────────────────────────────┐ │ │
│ │ │ 📋 规划:write_todos 分解任务 │ │ │
│ │ │ 📁 文件系统:读写中间结果 │ │ │
│ │ │ 🤖 子 Agent:task 委派 │ │ │
│ │ └────────────────────────────────────┘ │ │
│ └──────────────┬───────────────────────────┘ │
│ │ task 委派 │
│ ┌──────┼──────┐ │
│ ▼ ▼ ▼ │
│ [研究子Agent] [写作子Agent] [分析子Agent] │
│ 独立上下文 独立上下文 独立上下文 │
│ │ │ │ │
│ └──────┼──────┘ │
│ ▼ │
│ 共享文件系统交换结果 │
│ │ │
│ ▼ │
│ 主 Agent 更新计划,继续执行 │
│ │
└──────────────────────────────────────────────────┘
Deep Agent 的四大能力:
表格
python
from deepagents import create_deep_agent
from langchain_openai import ChatOpenAI
from langchain.tools import tool
from langgraph.checkpoint.memory import MemorySaver
# ---- 自定义工具 ----
@tool
def internet_search(query: str) -> str:
"""Run a web search."""
return f"Search results for: {query}"
# ---- 定义子 Agent ----
research_subagent = {
"name": "research",
"description": "Conducts focused web research on the topic",
"system_prompt": "You are a research assistant. Investigate the query thoroughly.",
"tools": [internet_search],
}
writer_subagent = {
"name": "writer",
"description": "Synthesizes information into a summary",
"system_prompt": "You are a skilled writer. Summarize key points clearly.",
"tools": [],
}
# ---- 创建 Deep Agent ----
agent = create_deep_agent(
model=ChatOpenAI(model="gpt-4o-mini", temperature=0),
tools=[internet_search],
system_prompt="""You are an expert research assistant.
Break the user's request into steps (using write_todos),
then gather information and write a final summary.""",
subagents=[research_subagent, writer_subagent],
checkpointer=MemorySaver(),
)
# ---- 运行 ----
result = agent.invoke({
"messages": [{"role": "user", "content": "Explain how rainbows form"}]
})
print(result["messages"][-1].content)
Deep Agent 执行流程:
规划:调用
write_todos,列出步骤("搜索彩虹原理"→"总结光学原理"→"撰写报告")委派:调用
task生成研究子 Agent,子 Agent 使用internet_search搜索存储:中间结果写入文件系统,主 Agent 上下文保持干净
整合:写作子 Agent 读取文件,生成最终报告
更新:主 Agent 标记任务完成,继续下一步
Deep Agent 本质上就是一个 LangGraph 图——你仍然可以用
.stream()、interrupt、checkpointer等所有 LangGraph 能力。
七、与其他框架对比
LangGraph vs CrewAI vs AutoGen
表格
选型决策框架
plaintext
你的需求是什么?
│
├─ 需要精确控制执行流、状态持久化、人机协同
│ └─ → LangGraph
│
├─ 快速原型、角色分工明确、不需要复杂状态管理
│ └─ → CrewAI
│
├─ 多专家对话/辩论/协商场景
│ └─ → AutoGen
│
└─ 需要子图嵌套、Map-Reduce、生产级可靠性
└─ → LangGraph(唯一选择)
一句话总结:LangGraph 是唯一原生支持子图嵌套和 Map-Reduce 的框架。如果你的系统复杂度到了需要这些能力的程度,别犹豫,选 LangGraph。
八、踩坑记录
坑 1:子图状态合并的双重拼接
前面已经详细说了。核心就一句话:共享键用 operator.add 时,子图返回的值会被 Reducer 再拼一次,导致重复。
python
# ❌ 裸 operator.add 在共享键上会双重拼接
class State(TypedDict):
path: Annotated[list[str], operator.add] # 会重复!
# ✅ 用自定义 Reducer 去重
class State(TypedDict):
path: Annotated[list[str], reduce_unique] # 不会重复
坑 2:子图内 interrupt 的传播
子图内部使用了 interrupt(),中断信号会传播到父图,但恢复时需要在正确的命名空间下操作。
python
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt
# 子图内有 interrupt
def human_review(state):
verdict = interrupt("请审核此结果") # ← 子图内的中断
return {"approved": verdict}
# 编译父图时必须传 checkpointer
graph = builder.compile(checkpointer=MemorySaver())
# 执行到 interrupt 会暂停
result = graph.invoke({"input": "..."}, config={"configurable": {"thread_id": "1"}})
# 恢复时需要传同一个 thread_id
result = graph.invoke(
{"approved": True}, # 提供中断时请求的值
config={"configurable": {"thread_id": "1"}},
)
注意:如果你在 LangGraph Studio 或 LangGraph API 中运行,子图的 interrupt 会自动被上层捕获,不需要手动处理命名空间。但本地运行时需要确保
thread_id一致。
坑 3:大量子图并行时的资源竞争
当你并行执行很多子图(比如 10+ 个 Send 实例),可能遇到:
LLM API 限流:10 个子图同时调用 API,瞬间触发 rate limit
内存压力:每个子图都有自己的状态副本,10 个并行 = 10 倍内存
解法:
python
# 方案1:控制并行度(LangGraph 0.2.x 支持)
graph = builder.compile()
result = graph.invoke(
input_data,
config={"recursion_limit": 50, "max_concurrency": 3}, # 最多 3 个并行
)
# 方案2:分批处理
from itertools import batched
subjects = ["lions", "cats", "dogs", "Python", "Rust", "Go", ...]
for batch in batched(subjects, 3): # 每 3 个一批
# 处理当前批次
pass
坑 4:子图嵌套层数不要超过 3 层
plaintext
❌ 不推荐:
Parent → Child → GrandChild → GreatGrandChild
↑ 第 4 层,调试困难,状态传递复杂
✅ 推荐:
Parent → Child → GrandChild
↑ 最多 3 层
理由:
调试困难:3 层嵌套时,trace 日志已经很长了,4 层基本无法阅读
状态传递损耗:每多一层,状态转换代码就多一份,出错概率指数级增长
性能损耗:每层子图编译和 checkpoint 都有开销
可维护性:没人愿意维护一个 4 层嵌套的状态机
如果确实需要更深的层次,考虑扁平化:把多层子图拆成一层 Supervisor 管理多个平级子图。
总结
表格
子图是 LangGraph 区别于其他多 Agent 框架的核心能力。它不只是代码组织方式的改进,更是一种架构思维:把复杂系统分解为可独立理解、独立测试、独立演进的模块。
从单图到子图,就像从单体应用到微服务的跨越——前期多花点设计时间,后期维护成本指数级下降。
📌 本文代码基于 langgraph 0.2.x,Python 3.11+
相关阅读:
第19篇 LangGraph StateGraph 状态机设计
第21篇 多Agent协作编排理论
评论区