LangGraph-Send-API并行任务与MapReduce模式
*LangGraph Send API MapReduce 模式——dispatcher 分发任务,Workers 并行执行,Reducer 合并结果*
LangGraph Send API:并行任务与 Map-Reduce 模式
1.1 静态并行的局限
LangGraph Send API MapReduce 模式——dispatcher 分发任务,Workers 并行执行,Reducer 合并结果
LangGraph 提供了一种声明并行的基础方式:在 add_edge 时同时指向多个节点,让它们同时执行。这种方式在任务数量固定时运作良好。
# 静态并行:编译时确定分支数量
builder.add_edge("start", "branch_a")
builder.add_edge("start", "branch_b")
builder.add_edge("start", "branch_c")
但现实中很多任务是动态的:用户传入 5 个 URL 要爬取,或者文档分块后产生 23 个片段需要逐一摘要,数量在运行前根本不知道。
RunnableParallel(LangChain 的并行执行工具,擅长固定数量的并发)同样面临这个问题。它擅长固定结构的并发,示例如下:
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 条消息",每条消息触发该节点的一次独立执行。
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 状态定义
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 阶段:派发任务
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 单文档处理节点
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 阶段:汇总节点
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 图的构建与连接
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 运行示例
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 流程图
1.5 实战:批量 URL 并行抓取与摘要生成
以下是一个更贴近生产的示例:并行抓取多个网页,对每个页面生成摘要,最终汇总。
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 阶段决定如何处理:
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 策略二:重试机制
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 层面设置超时:
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 的所有特性——持久化、流式输出、人工介入。