LangGraph 之 【工作流模式】(Send)
目录
1. 模式一:提示链模式(Prompt Chaining)
2. 模式二:并行化模式(Parallelization)
3. 模式三:路由模式(Routing)
4. 模式四:协调者-工作者模式(Orchestrator-Workers)
Send
5. 模式五:评估器-优化器模式(Evaluator-Optimizer)
6. 总结
1. 模式一:提示链模式(Prompt Chaining)
- 核心概念:这是最基础的线性模式,前一个节点的输出严格作为后一个节点的输入,就像写文章必须经历“大纲 → 初稿 → 润色 → 终稿”一样,每一步都不可跳跃
- 代码实现要点:通过定义 InputState(用户输入)和 OutputState(最终结果)来隔离内部状态,利用 OverallState 传递中间产物(大纲、初稿)
2. 模式二:并行化模式(Parallelization)
- 核心概念:多个独立任务同时执行,最后汇总结果。常用于多维度分析(如市场、竞品、技术同步调研)
- 代码实现要点:从 START 节点引出多条边,指向不同的分析节点,最后汇聚到汇总节点
3. 模式三:路由模式(Routing)
- 核心概念:根据输入内容动态选择分支。经典的智能客服场景:问题被分类为售前、售后或技术,交由不同的 Handler 处理
- 代码实现要点:利用 with_structured_output 强制 LLM 输出枚举值(Literal),配合 add_conditional_edges 进行分发
4. 模式四:协调者-工作者模式(Orchestrator-Workers)
- 核心概念:协调者(Orchestrator)在运行时分析任务,动态生成 n 个子任务,并通过 Send API 将不同数据分发到 n 个工作者节点并行执行
- 代码实现要点:assign_workers 函数返回一个 Send 列表,动态决定并发的数量
协调者-工作者模式和并行化模式都涉及同时执行多个任务,但它们的核心区别在于任务分配方式
- 并行化:任务在设计时就确定,所有任务同时开始
- 协调者-工作者:任务在运行时由协调者动态分配
Send
Send 的核心作用是让一个节点在运行时,
能动态地、并行地调用其他节点,并为每个调用传递定制化的状态
参数说明
| 参数 | 类型 | 是否必填 | 描述 |
|---|---|---|---|
node | str | 是 | 目标节点的名称,必须与使用add_node添加时的名称完全一致 |
arg | Any | 是 | 传递给目标节点的状态数据,可以是一个字典或自定义对象 |
timeout | float/timedelta/TimeoutPolicy | 否 | 为该特定任务单独设置的超时策略,会覆盖目标节点自身的超时配置 |
行为与机制
| 特性 | 说明 |
|---|---|
| 执行模式 | 从条件边返回多个Send对象,目标节点会被并行执行 |
| 状态更新 | 每个并行任务拥有独立的状态副本。任务完成后,可通过主 State 中定义的 Reducer(如operator.add)将结果合并回主状态 |
| 错误处理 | 可为每个任务单独设置超时时间。 某个任务失败,默认会影响整个图执行,需结合 Command等进行精细控制 |
import operator from typing import Annotated, TypedDict from langgraph.graph import StateGraph, START, END from langgraph.types import Send # 1. 定义主图状态 class OverallState(TypedDict): subjects: list[str] # 输入的主题列表 # jokes 字段使用 operator.add 作为 reducer,用于合并所有并发生成的笑话 jokes: Annotated[list[str], operator.add] # 2. 定义条件边函数 def continue_to_jokes(state: OverallState): """ 该函数从 START 节点被调用,它根据状态动态生成一个 Send 对象列表。 """ # 为 subjects 列表中的每一个主题,创建一个 Send 对象 # 每个 Send 对象都指向 "generate_joke" 节点,并传递一个包含该主题的字典 return [Send("generate_joke", {"subject": s}) for s in state["subjects"]] # 3. 定义工作节点 def generate_joke(state: dict): # 此节点接收由 Send 传来的定制状态 {"subject": s} subject = state["subject"] # 生成一个笑话,并返回一个字典,该字典会通过 reducer 合并到主状态 return {"jokes": [f"Joke about {subject}"]} # 4. 构建图 builder = StateGraph(OverallState) builder.add_node("generate_joke", generate_joke) # 关键:从 START 节点添加一个条件边,指向 continue_to_jokes 函数 builder.add_conditional_edges(START, continue_to_jokes) # 所有由 Send 触发的 "generate_joke" 节点执行完后,进入 END builder.add_edge("generate_joke", END) graph = builder.compile() # 5. 执行 result = graph.invoke({"subjects": ["cats", "dogs"]}) print(result) # 输出:{'subjects': ['cats', 'dogs'], 'jokes': ['Joke about cats', 'Joke about dogs']}- 并行局限性:并行任务不能直接修改主图的核心状态。必须通过返回一个字典,并依赖主 State 中定义的 Reducer 函数(如 operator.add)来合并结果
- 目标节点要求:被 Send 调用的节点(如示例中的 generate_joke),其函数签名接收的状态必须是 Send 传递的定制状态,而不是主图的完整状态
5. 模式五:评估器-优化器模式(Evaluator-Optimizer)
- 核心概念:先执行任务,再评估质量;若不达标,带着反馈重新执行
- 代码实现要点:循环结构:生成 → 评估 → (不合格)回到生成
6. 总结
| 业务场景 | 推荐模式 | 核心工程准则 |
|---|---|---|
| 内容生成、数据 ETL | 提示链 | 加入interrupt人工断点,做好 checkpoint 备份 |
| 独立多源数据同时检索 | 并行化 | 必须处理单个任务的异常降级,避免“一颗老鼠屎坏了一锅粥” |
| 客服、指令分类 | 路由 | 必须预留Fallback兜底节点,别信 LLM 会 100% 听话 |
| 长文档拆解、多章节写作 | 协调者-工作者 | 必须在汇总时按Index重排序,并控制并发数 |
| 代码修复、高质量翻译 | 评估器-优化器 | 必须设置循环上限(max_trials),防止死循环烧钱 |
