课程0基础Agent开发课 / LangGraph / LangGraph-Send-API并行任务与MapReduce模式
— 20 min read

LangGraph-Send-API并行任务与MapReduce模式

*LangGraph Send API MapReduce 模式——dispatcher 分发任务,Workers 并行执行,Reducer 合并结果*

LangGraph Send API:并行任务与 Map-Reduce 模式


1.1 静态并行的局限

输入任务

Dispatcher
分发任务

Worker 1
并行执行

Worker 2
并行执行

Worker N
并行执行

Reducer
合并结果

输出结果

LangGraph Send API MapReduce 模式——dispatcher 分发任务,Workers 并行执行,Reducer 合并结果

LangGraph 提供了一种声明并行的基础方式:在 add_edge 时同时指向多个节点,让它们同时执行。这种方式在任务数量固定时运作良好。

python
# 静态并行:编译时确定分支数量
builder.add_edge("start", "branch_a")
builder.add_edge("start", "branch_b")
builder.add_edge("start", "branch_c")

但现实中很多任务是动态的:用户传入 5 个 URL 要爬取,或者文档分块后产生 23 个片段需要逐一摘要,数量在运行前根本不知道。

RunnableParallel(LangChain 的并行执行工具,擅长固定数量的并发)同样面临这个问题。它擅长固定结构的并发,示例如下:

python
from langchain_core.runnables import RunnableParallel

chain = RunnableParallel(
    summary=summarize_chain,
    keywords=keyword_chain,
    sentiment=sentiment_chain,
)

这里有三个固定的并行分支,逻辑清晰。但如果要并行处理一个长度不定的列表,RunnableParallel 就无能为力了——它的结构必须在代码编写时确定,无法在运行时动态生成分支。

对比维度 RunnableParallel LangGraph Send API
并行数量 编译时固定 运行时动态生成
状态管理 无持久化状态 完整 State 机制
错误隔离 单个失败影响整体 可逐分支处理异常
检查点 不支持 支持断点续跑
适用场景 固定多路 LLM 调用 动态批处理、Map-Reduce

1.2 Send API 是什么

Send 是 LangGraph 在 0.1.x 版本引入的机制,允许在节点函数内部,在运行时动态地向某个节点派发任意数量的任务,每个任务携带独立的输入状态。

核心思想:节点不再只能"跳转到某个节点",而是可以"向某个节点发送 N 条消息",每条消息触发该节点的一次独立执行。

python
from langgraph.types import Send

def generate_tasks(state: State) -> list[Send]:
    return [Send("process_item", {"item": x}) for x in state["items"]]

这个函数的返回值不是下一个节点的名字,而是一个 Send 对象列表。每个 Send 对象指定:

  • 第一个参数:目标节点名称
  • 第二个参数:传递给该节点的输入状态

LangGraph 运行时会并发执行所有 Send 派发的任务,并在所有任务完成后将结果合并到主状态中。


1.3 Map-Reduce 完整实现

Map-Reduce(映射-归约:一种分布式计算模式,Map 阶段把一批任务拆开并行处理,Reduce 阶段把所有结果汇聚合并,最终得到统一输出)是 Send API 最典型的应用场景:Map 阶段并行处理每个元素,Reduce 阶段汇总所有结果。

以下示例实现一个文档批量摘要系统:输入多个文档,并行生成每份摘要,最后合并成一份总结报告。

1.3.1 状态定义

python
from typing import Annotated, TypedDict
import operator
from langchain_core.documents import Document

class DocumentState(TypedDict):
    """主状态:贯穿整个图"""
    documents: list[Document]      # 待处理的文档列表
    summaries: Annotated[list[str], operator.add]  # 汇总器:自动追加
    final_report: str              # 最终报告

class SingleDocState(TypedDict):
    """单文档处理节点的独立状态"""
    document: Document
    summaries: Annotated[list[str], operator.add]

Annotated[list[str], operator.add] 是关键:它告诉 LangGraph,当多个并行分支都向 summaries 写入数据时,不要用覆盖的方式,而是用 operator.add(列表拼接)来合并。没有这个注解,并行写入会互相覆盖,只保留最后一个结果。

1.3.2 Map 阶段:派发任务

python
from langgraph.types import Send

def dispatch_documents(state: DocumentState) -> list[Send]:
    """Map 阶段:为每个文档创建一个独立的处理任务"""
    return [
        Send("summarize_document", {"document": doc, "summaries": []})
        for doc in state["documents"]
    ]

注意这个函数没有返回节点名,而是返回 list[Send]。LangGraph 会将这个函数注册为条件边(conditional edge),并识别返回值是 Send 列表。

1.3.3 单文档处理节点

python
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage

llm = ChatOpenAI(model="gpt-4o-mini", temperature=0)

def summarize_document(state: SingleDocState) -> dict:
    """处理单个文档,生成摘要"""
    doc = state["document"]

    response = llm.invoke([
        HumanMessage(content=f"""请为以下文档生成一段简洁的摘要(100字以内):

文档标题:{doc.metadata.get('title', '无标题')}
文档内容:{doc.page_content[:2000]}

摘要:""")
    ])

    return {"summaries": [response.content]}

每次执行这个节点时,它只处理一个文档,返回值中的 summaries 列表只包含当前文档的摘要。operator.add reducer 会负责把所有并行实例的结果拼接在一起。

1.3.4 Reduce 阶段:汇总节点

python
def generate_report(state: DocumentState) -> dict:
    """Reduce 阶段:汇总所有摘要,生成综合报告"""
    summaries_text = "\n\n".join([
        f"文档 {i+1}{summary}"
        for i, summary in enumerate(state["summaries"])
    ])

    response = llm.invoke([
        HumanMessage(content=f"""以下是多份文档的摘要,请综合生成一份总结报告:

{summaries_text}

综合报告:""")
    ])

    return {"final_report": response.content}

1.3.5 图的构建与连接

python
from langgraph.graph import StateGraph, START, END

builder = StateGraph(DocumentState)

# 注册节点
builder.add_node("summarize_document", summarize_document)
builder.add_node("generate_report", generate_report)

# START -> dispatch_documents(作为条件边,返回 Send 列表)
builder.add_conditional_edges(
    START,
    dispatch_documents,
    ["summarize_document"],  # 声明可能的目标节点
)

# 所有 summarize_document 完成后 -> generate_report
builder.add_edge("summarize_document", "generate_report")
builder.add_edge("generate_report", END)

graph = builder.compile()

1.3.6 运行示例

python
from langchain_core.documents import Document

documents = [
    Document(
        page_content="LangGraph 是一个基于图的 Agent 框架,支持状态持久化和复杂流程控制...",
        metadata={"title": "LangGraph 简介"}
    ),
    Document(
        page_content="向量数据库是 RAG 系统的核心组件,常用方案包括 Chroma、Milvus、Pinecone...",
        metadata={"title": "向量数据库概述"}
    ),
    Document(
        page_content="Prompt Engineering 是提升 LLM 输出质量的关键技术,包括 Few-shot、Chain-of-Thought...",
        metadata={"title": "Prompt 工程指南"}
    ),
]

result = graph.invoke({"documents": documents, "summaries": []})
print(result["final_report"])

1.4 Map-Reduce 流程图

Send doc_1

Send doc_2

Send doc_N

summaries += [摘要1]

summaries += [摘要2]

summaries += [摘要N]

开始

\dispatch_documents
(条件边:返回 Send 列表)\

\summarize_document
实例1:处理文档1\

\summarize_document
实例2:处理文档2\

\summarize_document
实例N:处理文档N\

Annotated Reducer
operator.add 合并

\generate_report
汇总所有摘要\

结束


1.5 实战:批量 URL 并行抓取与摘要生成

以下是一个更贴近生产的示例:并行抓取多个网页,对每个页面生成摘要,最终汇总。

python
import asyncio
import httpx
from bs4 import BeautifulSoup
from typing import Annotated, TypedDict
import operator
from langgraph.types import Send
from langgraph.graph import StateGraph, START, END
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage

# 状态定义
class CrawlState(TypedDict):
    urls: list[str]
    page_summaries: Annotated[list[dict], operator.add]
    final_digest: str

class SingleUrlState(TypedDict):
    url: str
    page_summaries: Annotated[list[dict], operator.add]

llm = ChatOpenAI(model="gpt-4o-mini", temperature=0)

# 网页抓取函数
def fetch_page_content(url: str) -> str:
    """同步抓取网页内容"""
    try:
        with httpx.Client(timeout=10.0, follow_redirects=True) as client:
            response = client.get(url, headers={
                "User-Agent": "Mozilla/5.0 (compatible; ResearchBot/1.0)"
            })
            soup = BeautifulSoup(response.text, "html.parser")
            # 移除脚本和样式标签
            for tag in soup(["script", "style", "nav", "footer"]):
                tag.decompose()
            return soup.get_text(separator="\n", strip=True)[:3000]
    except Exception as e:
        return f"抓取失败:{str(e)}"

# Map 阶段:派发 URL 任务
def dispatch_urls(state: CrawlState) -> list[Send]:
    return [
        Send("process_url", {"url": url, "page_summaries": []})
        for url in state["urls"]
    ]

# 单 URL 处理节点
def process_url(state: SingleUrlState) -> dict:
    url = state["url"]
    content = fetch_page_content(url)

    if content.startswith("抓取失败"):
        return {"page_summaries": [{"url": url, "summary": content, "status": "error"}]}

    response = llm.invoke([
        HumanMessage(content=f"请为以下网页内容生成摘要(150字以内):\n\n{content}\n\n摘要:")
    ])

    return {
        "page_summaries": [{
            "url": url,
            "summary": response.content,
            "status": "success"
        }]
    }

# Reduce 阶段:生成综合摘要
def compile_digest(state: CrawlState) -> dict:
    successful = [s for s in state["page_summaries"] if s["status"] == "success"]
    failed = [s for s in state["page_summaries"] if s["status"] == "error"]

    if not successful:
        return {"final_digest": "所有 URL 抓取失败,无法生成摘要。"}

    summaries_text = "\n\n".join([
        f"来源:{s['url']}\n摘要:{s['summary']}"
        for s in successful
    ])

    response = llm.invoke([
        HumanMessage(content=f"综合以下多个网页的摘要,生成一份简洁的信息汇总报告:\n\n{summaries_text}")
    ])

    report = response.content
    if failed:
        report += f"\n\n注:以下 {len(failed)} 个 URL 抓取失败:" + \
                  ", ".join(s["url"] for s in failed)

    return {"final_digest": report}

# 构建图
builder = StateGraph(CrawlState)
builder.add_node("process_url", process_url)
builder.add_node("compile_digest", compile_digest)

builder.add_conditional_edges(START, dispatch_urls, ["process_url"])
builder.add_edge("process_url", "compile_digest")
builder.add_edge("compile_digest", END)

graph = builder.compile()

# 使用示例
urls = [
    "https://python.org",
    "https://docs.langchain.com",
    "https://openai.com/blog",
]

result = graph.invoke({"urls": urls, "page_summaries": []})
print(result["final_digest"])

1.6 错误处理:并行分支失败的策略

并行任务中,单个分支失败是常见情况(网络超时、内容解析错误等)。有三种处理策略:

1.6.1 策略一:局部异常捕获(推荐)

在每个处理节点内部捕获异常,返回带错误标记的结果,让 Reduce 阶段决定如何处理:

python
def process_url(state: SingleUrlState) -> dict:
    url = state["url"]
    try:
        content = fetch_page_content(url)
        summary = llm.invoke([HumanMessage(content=f"摘要:{content}")]).content
        return {"page_summaries": [{"url": url, "summary": summary, "status": "ok"}]}
    except Exception as e:
        # 捕获异常,不抛出,返回错误标记
        return {"page_summaries": [{"url": url, "summary": "", "status": "error", "error": str(e)}]}

1.6.2 策略二:重试机制

python
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 fetch_with_retry(url: str) -> str:
    return fetch_page_content(url)

def process_url_with_retry(state: SingleUrlState) -> dict:
    url = state["url"]
    try:
        content = fetch_with_retry(url)
        summary = llm.invoke([HumanMessage(content=f"摘要:{content}")]).content
        return {"page_summaries": [{"url": url, "summary": summary, "status": "ok"}]}
    except Exception as e:
        return {"page_summaries": [{"url": url, "summary": "", "status": "failed"}]}

1.6.3 策略三:超时控制

对于外部 IO 操作,建议在 graph.invoke 层面设置超时:

python
from langgraph.checkpoint.memory import MemorySaver

memory = MemorySaver()
graph = builder.compile(checkpointer=memory)

config = {
    "configurable": {"thread_id": "batch_001"},
    "recursion_limit": 100,
}

result = graph.invoke(
    {"urls": urls, "page_summaries": []},
    config=config,
)

三种策略的适用场景:

策略 适用场景 代价
局部捕获 允许部分失败,Reduce 阶段可以处理缺失数据 最低,推荐默认选择
重试机制 失败原因是瞬时的(网络抖动、限流) 增加执行时间
超时控制 有严格的 SLA 要求 可能丢失部分结果

1.7 Send API 与 RunnableParallel 的选型原则

明确两者的定位差异,可以避免选型错误:

选 RunnableParallel 的场景:

  • 固定数量的并行 LLM 调用(如同时生成摘要和关键词)
  • 不需要 LangGraph 的状态管理
  • 简单的链式调用,不涉及复杂流程控制

选 Send API 的场景:

  • 运行时才知道任务数量(处理用户上传的文件列表、爬取 URL 列表)
  • 需要 Map-Reduce 模式处理批量数据
  • 需要与 LangGraph 的检查点、持久化特性结合
  • 单个任务失败不能影响其他任务的执行

Send API 做的是:把"循环处理列表"转换为图的原生并行执行。原本 for 循环串行的逻辑,改用 Send API 之后并发执行,同时保留 LangGraph 的所有特性——持久化、流式输出、人工介入。

本页目录