目 录CONTENT

文章目录

LangGraph 子图与多Agent编排:从模块化到协作模式

PySuper
2025-12-13 / 0 评论 / 0 点赞 / 0 阅读 / 0 字
温馨提示:
本文最后更新于2026-05-22,若内容或图片失效,请留言反馈。 所有牛逼的人都有一段苦逼的岁月。 但是你只要像SB一样去坚持,终将牛逼!!! ✊✊✊

作者: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"}              │
│                                                          │
└──────────────────────────────────────────────────────────┘

状态键映射规则

共享键模式下,核心规则只有三条:

  1. 子图的 State 必须包含所有与父图共享的键,子图可以有额外私有键

  2. 共享键上的 Reducer 必须一致(父子图的同一个键用同一个 Reducer)

  3. 子图返回时,只有共享键的值会传回父图,私有键留在子图内部

三、子图的状态隔离

私有状态 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()

注意:failuressummary 是子图私有状态,父图完全不知道它们的存在。这就是状态隔离的力量。

⚠️ 大坑: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 的关键点

表格

要点

说明

Reducer 必须支持合并

多个 Send 实例的结果需要合并,对应字段必须用 operator.add 或自定义 Reducer

每个 Send 独立执行

彼此之间看不到对方的状态

路由函数返回 list[Send]

不是返回节点名字符串,而是 Send 对象列表

动态数量

主题数量在运行时才确定,图定义时不需要知道

六、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 的四大能力:

表格

能力

工具

作用

规划

write_todos

分解复杂任务为步骤

文件系统

ls, read_file, write_file, edit_file

上下文卸载,防止窗口溢出

子Agent委派

task

生成专门子Agent处理子任务,隔离上下文

长期记忆

LangGraph Store

跨会话持久化信息

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 执行流程:

  1. 规划:调用 write_todos,列出步骤("搜索彩虹原理"→"总结光学原理"→"撰写报告")

  2. 委派:调用 task 生成研究子 Agent,子 Agent 使用 internet_search 搜索

  3. 存储:中间结果写入文件系统,主 Agent 上下文保持干净

  4. 整合:写作子 Agent 读取文件,生成最终报告

  5. 更新:主 Agent 标记任务完成,继续下一步

Deep Agent 本质上就是一个 LangGraph 图——你仍然可以用 .stream()interruptcheckpointer 等所有 LangGraph 能力。

七、与其他框架对比

LangGraph vs CrewAI vs AutoGen

表格

维度

LangGraph

CrewAI

AutoGen

核心范式

图状态机

角色扮演团队

对话驱动

学习曲线

陡峭(图论 + 状态管理)

平缓(声明式角色定义)

中等(事件驱动)

控制粒度

节点级,最细粒度

任务级

消息级

状态管理

显式 TypedDict + Reducer

隐式,框架管理

隐式,对话历史

Token 效率

⭐⭐⭐⭐⭐ 状态隔离减少上下文

⭐⭐⭐ 全量共享

⭐⭐ 对话全量共享

持久化

原生 Checkpointer

基础 Memory

基础

人机协同

interrupt 节点级

任务级 human_input

human_input_mode 消息级

子图/嵌套

✅ 原生支持

❌ 无子图概念

❌ 无子图概念

并行执行

✅ 原生 fan-out/fan-in

✅ 并行任务

⚠️ 有限

Map-Reduce

✅ Send API

生态集成

LangChain 700+ 集成

官方工具市场 + MCP

社区工具库

生产部署

LangGraph Cloud + Studio

CrewAI Enterprise

AutoGen Studio

语言支持

Python + JS/TS

Python

Python + .NET

选型决策框架

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 层

理由:

  1. 调试困难:3 层嵌套时,trace 日志已经很长了,4 层基本无法阅读

  2. 状态传递损耗:每多一层,状态转换代码就多一份,出错概率指数级增长

  3. 性能损耗:每层子图编译和 checkpoint 都有开销

  4. 可维护性:没人愿意维护一个 4 层嵌套的状态机

如果确实需要更深的层次,考虑扁平化:把多层子图拆成一层 Supervisor 管理多个平级子图。

总结

表格

主题

核心收获

子图价值

模块化、复用、团队分工、状态隔离

两种接入方式

共享键直接 add_node vs 包装函数状态转换

状态隔离

父图只看输入输出,私有状态不可见

operator.add 大坑

共享键双重拼接,用自定义 Reducer 或包装函数解决

Supervisor 模式

最常用的多 Agent 编排,一个总控 + N 个专家

并行召回

fan-out/fan-in + Reducer 合并去重

Map-Reduce

Send API 实现动态数量并行,Reducer 自动合并

Deep Agents

规划 + 文件系统 + 子Agent 委派,处理复杂长任务

框架选型

需要子图/Map-Reduce/精确控制 → LangGraph,其他场景看需求

嵌套层数

最多 3 层,超过就扁平化

子图是 LangGraph 区别于其他多 Agent 框架的核心能力。它不只是代码组织方式的改进,更是一种架构思维:把复杂系统分解为可独立理解、独立测试、独立演进的模块。

从单图到子图,就像从单体应用到微服务的跨越——前期多花点设计时间,后期维护成本指数级下降。

📌 本文代码基于 langgraph 0.2.x,Python 3.11+

相关阅读:

0
  1. 支付宝打赏

    qrcode alipay
  2. 微信打赏

    qrcode weixin

评论区