愿旖旎头像
关注
【LangGraph实战】LangGraph 学习笔记(三):Overwrite、输入输出模式与四大工作流模式封面图

【LangGraph实战】LangGraph 学习笔记(三):Overwrite、输入输出模式与四大工作流模式

   👋 欢迎阅读

  🔥 愿旖旎 · 个人主页

📘 学习专栏: 《算法专栏》《LangChain学习》《贪心算法》

🌄 钱塘江上潮信来,今日方知我是我

✨当前学习内容:《LangGraph》

📑 目录


一、前置讲解

在进入正文前,先花 10 秒了解本篇会反复用到的核心概念,带着印象去读正文,学习效率更高。

🧠 前置知识一:Reducer 与 Overwrite

Reducer 决定状态怎么合并(追加还是覆盖),而 Overwrite 能强制绕过它,直接整体替换。

🧩 前置知识二:input_schema 与 output_schema

让图只收该收的输入、只吐该吐的输出,内部状态可以比对外接口更复杂。

🔒 前置知识三:私有状态(Private State)

用类型注解给节点开一条私有数据通道:节点之间能传,但不出现在最终输出里。

🔀 前置知识四:扇出与扇入(Fan-out / Fan-in)

START 同时指向多个节点就是扇出(并行),多个节点汇合到一个节点就是扇入。

📮 前置知识五:Send(动态派发任务)

在条件边里返回 Send("节点名", 状态),让并行分支的数量在运行时才决定。


二、LangGraph 的三个进阶特性

前面两篇我们用 State、Nodes、Edges、Reducer 搭出了搜索 Agent 和代理式 RAG。这一篇先补三个"进阶开关",它们直接决定了状态怎么被改写、图的对外接口长什么样、数据谁可见。

2.1 使用 Overwrite 绕过 reducer

在 LangGraph 中,reducers 用于控制状态更新的处理方式。默认情况下,每个状态键都有独立的 reducer 函数,决定如何合并节点返回的更新。但有时候我们需要完全覆盖状态值而不是合并,这时就要用 Overwrite。

最典型的场景是聊天应用:

情况期望行为靠什么实现
正常对话新消息追加到消息列表默认的 operator.add reducer
清空聊天记录重新开始直接覆盖整个消息列表Overwrite 强制绕过 reducer

代码示例:

from langgraph.graph import StateGraph, START, END
from langgraph.types import Overwrite
from typing_extensions import Annotated, TypedDict
import operator
​
​
# ---------- 1. 定义状态(State)----------
# State 是图中流转的"共享数据",用 TypedDict 声明字段结构
class State(TypedDict):
    # Annotated[list, operator.add] 中的 operator.add 是【reducer(归约函数)】
    # 作用:节点返回的新值会与旧值"相加/追加",而不是直接覆盖
    messages: Annotated[list, operator.add]
​
​
# ---------- 2. 定义节点函数(每个节点就是一个普通函数)----------
def add_message(state: State):
    """节点1:往 messages 追加一条消息(走 reducer,会追加)"""
    return {"messages": ["first message"]}
​
​
def replace_messages(state: State):
    """节点2:用 Overwrite 绕过 reducer,直接【替换】整个列表"""
    # Overwrite(...):强制绕过 reducer,不做追加,而是整体替换
    return {"messages": Overwrite(["replacement message"])}
​
​
# ---------- 3. 构建图(Graph)----------
builder = StateGraph(State)          # 创建状态图构建器(绑定 State)
​
# add_node:添加【节点】(名称, 函数)—— 节点 = 图中执行的一步
builder.add_node("add_message", add_message)
builder.add_node("replace_messages", replace_messages)
​
# add_edge:添加【边】—— 定义节点的执行顺序(从一个节点走到另一个)
builder.add_edge(START, "add_message")                    # 入口:START → add_message
builder.add_edge("add_message", "replace_messages")       # add_message → replace_messages
builder.add_edge("replace_messages", END)                 # 出口:replace_messages → END
​
# compile:编译图,生成可执行的对象
graph = builder.compile()
​
​
# ---------- 4. 执行图 ----------
# invoke:传入初始状态,运行整张图
result = graph.invoke({"messages": ["initial"]})
print(result["messages"])

运行结果:

💡 对照看这两个结果就很清楚了:不加 Overwrite 时,replace_messages 返回的列表被 reducer 当成"要追加的内容",所以三条消息全在;加上 Overwrite 后,整个列表被整体替换,连最初的 initial 也一起没了。

使用 Overwrite 的应用场景:

场景说明
重置对话清空聊天历史,开始新对话
状态重置在错误恢复后重置应用状态
数据清理替换损坏或过时的数据

2.2 定义输入输出模式

默认情况下,LangGraph 使用单一的状态模式。但我们可以定义独立的输入和输出模式:

模式作用
独立的输入模式验证输入数据的结构
独立的输出模式过滤输出数据,只返回需要的信息
内部模式节点间通信使用的完整状态

为什么需要独立的内部模式? 考虑一个问答系统:

环节内容
输入用户的问题(字符串)
输出AI 的答案(字符串)
内部可能需要存储中间结果、上下文等信息

对应到 StateGraph 的三个参数:

参数名类型描述
input_schematype[InputT] 或 None定义 StateGraph 输入的 State 类
output_schematype[OutputT] 或 None定义 StateGraph 输出的 State 类
state_schematype[StateT]定义 StateGraph 内部的 State 类(第一个位置参数)

代码示例:

from langgraph.graph import StateGraph, START, END
from typing_extensions import TypedDict
​
​
# ---------- 1. 定义输入模式(只包含用户问题)----------
# input_schema:调用 invoke 时传入的状态,只允许有 question 字段
class InputState(TypedDict):
    question: str
​
​
# ---------- 2. 定义输出模式(只包含 AI 答案)----------
# output_schema:图跑完后返回的状态,只保留 answer 字段
class OutputState(TypedDict):
    answer: str
​
​
# ---------- 3. 定义完整状态模式(图内部实际使用)----------
# OverallState 同时继承 InputState 和 OutputState → 内部既有 question 又有 answer
class OverallState(InputState, OutputState):
    pass
​
​
# ---------- 4. 定义节点函数 ----------
def answer_node(state: InputState):
    """处理输入并生成答案"""
    # 这里可以访问 question,生成 answer
    # 返回两个字段:answer(新生成)+ question(原样带回,供内部流转)
    return {
        "answer": f"Answer to: {state['question']}",
        "question": state["question"],
    }
​
​
# ---------- 5. 构建图时指定输入 / 输出模式 ----------
builder = StateGraph(
    OverallState,
    input_schema=InputState,    # 输入验证:只接受 question
    output_schema=OutputState,  # 输出过滤:只返回 answer
)
​
# add_node:添加节点;add_edge:添加边(定义执行顺序)
builder.add_node("answer_node", answer_node)
builder.add_edge(START, "answer_node")   # 入口 → answer_node
builder.add_edge("answer_node", END)     # answer_node → 出口
​
# compile:编译图
graph = builder.compile()
​
​
# ---------- 6. 测试 ----------
result = graph.invoke({"question": "What is LangGraph?"})
print(result)

运行结果:

{'answer': 'Answer to: What is LangGraph?'}

💡 节点内部明明返回了 question 和 answer 两个字段,但 print(result) 里只有 answer——这就是 output_schema 在过滤输出。同时 input_schema 还顺手做了输入校验:如果 invoke 时传了一个 InputState 里没有的键,会直接报错,而不是悄悄带进流程里。

独立的输入和输出模式,实际应用场景:

场景说明
API 开发定义清晰的请求/响应格式
微服务服务间明确的数据契约
数据管道明确的输入输出规范

2.3 在节点间传递私有状态

有时候节点之间需要传递临时数据:这些数据对中间逻辑很重要,但不该出现在最终输出里,只在特定节点间共享——这就是私有状态。

举个例子,要根据数据库里的信息生成数据报告,可以设置三个节点:

节点职责能看到什么
节点 1从数据库获取原始数据(含敏感信息)完整原始数据
节点 2处理数据,过滤掉敏感信息完整原始数据(要处理它)
节点 3生成最终报告只有清理后的结果

也就是说:节点 1 和节点 2 需要共享原始数据,但节点 3 不应该看到敏感信息。

代码示例:

from langgraph.graph import StateGraph, START, END
from typing_extensions import TypedDict
​
​
# ---------- 1. 定义公共状态(最终输出中可见)----------
class OverallState(TypedDict):
    final_result: str
​
​
# ---------- 2. 定义节点1的私有输出(不进入最终状态)----------
class Node1Output(TypedDict):
    sensitive_data: str      # 这个字段不会出现在最终状态中
​
​
# ---------- 3. 定义节点2需要的输入(包含私有数据)----------
class Node2Input(TypedDict):
    sensitive_data: str
​
​
# ---------- 4. 定义三个节点函数 ----------
# 每个节点用【类型注解】声明自己"需要什么输入"和"产出什么输出"
def node_1(state: OverallState) -> Node1Output:
    """第一步:获取包含敏感信息的原始数据"""
    private_data = "这是敏感信息"
    print("Node1: 获取到敏感数据,但不会暴露给最终输出")
    return {"sensitive_data": private_data}
​
​
def node_2(state: Node2Input) -> OverallState:
    """第二步:处理数据,移除敏感信息"""
    print(f"Node2: 处理敏感数据: {state['sensitive_data']}")
    # 处理数据,返回清理后的结果(只返回公共状态字段)
    return {"final_result": "清理后的处理结果"}
​
​
def node_3(state: OverallState) -> OverallState:
    """第三步:只看到清理后的数据"""
    print(f"Node3: 只能看到最终结果: {state['final_result']}")
    return {"final_result": state["final_result"] + " - 完成"}
​
​
# ---------- 5. 构建图 ----------
builder = StateGraph(OverallState)
​
# add_sequence:支持一次添加一系列节点,按所给顺序依次执行
# 注意:这里使用 add_sequence,但类型系统会处理私有状态
builder.add_sequence([node_1, node_2, node_3])
​
# add_edge:添加入口边(START → node_1)
builder.add_edge(START, "node_1")
​
# compile:编译图
graph = builder.compile()
​
​
# ---------- 6. 测试 ----------
response = graph.invoke({"final_result": "initial"})
print(f"\n最终输出: {response}")

运行结果:

Node1: 获取到敏感数据,但不会暴露给最终输出
Node2: 处理敏感数据: 这是敏感信息
Node3: 只能看到最终结果: 清理后的处理结果
​
最终输出: {'final_result': '清理后的处理结果 - 完成'}

💡 注意 add_sequence([node_1, node_2, node_3]) 只是把三个节点的执行顺序串起来(内部自动补节点间的边),入口边仍要自己加:builder.add_edge(START, "node_1")。另外它取节点名的方式和 add_node(函数) 一样,默认用函数名。

Node1Output 和 Node2Input 到底有什么用?

这两个类的字段完全一样(都是 sensitive_data: str),看着像重复,其实它们定义的是节点之间的数据通道:

node_1 的产出是 Node1Output
      ↓ (同名同类型字段 sensitive_data 直接对接)
node_2 的输入是 Node2Input

也就是说:node_1 交出的 sensitive_data,正好被 node_2 读走,数据流是连通的。关键点在于:这类字段不在 OverallState(公共状态)里,所以它的可见性是"半开"的:

行为结果
节点之间传递(node_1 → node_2)✅ 能读到 sensitive_data
出现在最终输出里❌ 被过滤掉了

带来两个实际好处:

  • 数据可见性可控:敏感信息、API Key、中间大对象只在节点间流转,不对外暴露;

  • 对外只返回干净字段:外部拿到的就是干净的 final_result。

能不能合并成一个类? 能!因为字段一样,写成同一个类完全等价:

# 合并写法(与拆成两个类等价)
class SensitiveData(TypedDict):
    sensitive_data: str
​
​
def node_1(state: OverallState) -> SensitiveData:
    ...
​
​
def node_2(state: SensitiveData) -> OverallState:
    ...

💡 教程里之所以写成两个类,是为了让"谁的输出 / 谁的输入"这层语义更清晰;真正跑起来,一个类也完全没问题。

节点间传递私有状态,实际应用场景:

场景说明
数据处理中间处理步骤的临时数据
认证流程令牌等敏感信息的传递
复杂计算中间计算结果
错误处理错误详情在内部传递,但对外提供友好消息

以上三个特性让 LangGraph 能够处理复杂的企业级应用场景,同时保持代码的清晰和安全性。


三、工作流的四大常见模式

工作流模式是预先定义好的执行路径,就像工厂的流水线一样,每个步骤都有明确的输入输出和顺序。不同需求场景,对应不同的模式选择。

模式一句话概括典型场景
提示链一条直线,前一步的输出喂给后一步大纲 → 初稿 → 润色 → 终稿
并行化多个任务同时干,最后汇总多维度分析、批量处理
路由先分类,再分流到专用处理器智能客服、工单分派
协调者-工作者一个大脑拆任务,多个工人并行干长文档分章节撰写

3.1 提示链模式(Prompt Chaining)

3.1.1 概念

提示链就像流水线一样,前一个步骤的输出作为下一个步骤的输入。这跟做内容创作时一样:需要 大纲 → 初稿 → 润色 → 最终稿,每一步的输出都要传给下一步,才能让内容质量逐步提升。

📷 插图位置:提示链模式示意图(START → 输入主题 → 大纲生成 → 大纲 → 初稿生成 → 初稿 → 润色文章 → 润色版 → 最终稿生成 → 终稿 → END)

3.1.2 模式实践

我们创建一个内容创作场景的工作流,包含 大纲 → 初稿 → 润色 → 最终稿。节点设计为:

节点职责
generate_outline只负责大纲生成
generate_draft只负责初稿写作
polish_content只负责内容润色
finalize_content只负责最终整合

代码示例:

from langchain.chat_models import init_chat_model
from langchain_core.messages import HumanMessage
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END

# 定义大模型(适配 DeepSeek)
model = init_chat_model(model="deepseek-chat", model_provider="deepseek")


# ---------- 1. 定义输入模式(只包含用户输入)----------
class InputState(TypedDict):
   topic: str              # 用户输入的主题


# ---------- 2. 定义输出模式(只包含最终结果)----------
class OutputState(TypedDict):
   final_content: str      # 最终的内容


# ---------- 3. 定义完整状态模式(内部使用)----------
class OverallState(InputState, OutputState):
   outline: str            # 第一步:生成的大纲
   draft: str              # 第二步:生成的初稿
   polished_draft: str     # 第三步:润色后的稿件


# ---------- 4. 第一步:生成大纲 ----------
PROMPT_1 = (
   "根据主题生成文章大纲。\n"
   "主题:{topic}\n"
   "要求:"
   "1.只需两个最核心标题"
   "2.不用进行说明,只返回最终大纲"
)


def generate_outline(state: InputState):
   """根据主题生成内容大纲"""
   print("*" * 50)
   print(f"内容大纲生成中...\n")
   prompt = PROMPT_1.format(topic=state['topic'])
   outline = model.invoke([HumanMessage(content=prompt)]).content
   print(f"大纲已生成:\n{outline}\n")
   return {
       "outline": outline,
       "topic": state["topic"],
   }


# ---------- 5. 第二步:基于大纲生成初稿 ----------
PROMPT_2 = (
   "根据以下内容生成文章完整初稿。\n"
   "主题:{topic}\n"
   "大纲: "
   "{outline}\n"
   "要求:"
   "1.每个标题下,最多使用三句话的内容即可"
   "2.不用进行说明,只返回最终结果"
)


def generate_draft(state: OverallState):
   """根据大纲生成完整初稿"""
   print("*" * 50)
   print(f"生成初稿中...\n")
   prompt = PROMPT_2.format(topic=state['topic'], outline=state['outline'])
   draft = model.invoke([HumanMessage(content=prompt)]).content
   print(f"初稿已生成:\n{draft}\n")
   return {"draft": draft}


# ---------- 6. 第三步:润色稿件 ----------
PROMPT_3 = (
   "根据文章初稿进行润色。\n"
   "主题:{topic}\n"
   "初稿: "
   "{draft}\n"
   "要求:"
   "1.润色后,文章不能太长"
)


def polish_content(state: OverallState):
   """对初稿进行润色优化"""
   print("*" * 50)
   print(f"文章润色中...\n")
   prompt = PROMPT_3.format(topic=state['topic'], draft=state['draft'])
   polished = model.invoke([HumanMessage(content=prompt)]).content
   print(f"润色完成,内容如下:\n{polished}\n")
   return {"polished_draft": polished}


# ---------- 7. 第四步:生成最终稿 ----------
PROMPT_4 = (
   "根据润色版文章,生成文章终稿。\n"
   "主题:{topic}\n"
   "大纲: "
   "{outline}\n"
   "润色版文章: "
   "{polished_draft}\n"
)


def finalize_content(state: OverallState):
   """生成最终版本的内容"""
   prompt = PROMPT_4.format(
       topic=state['topic'],
       outline=state['outline'],
       polished_draft=state['polished_draft'],
   )
   final_content = model.invoke([HumanMessage(content=prompt)]).content
   return {"final_content": final_content}


# ---------- 8. 构建工作流 ----------
builder = StateGraph(
   OverallState,
   input_schema=InputState,     # 输入验证:只接受 topic
   output_schema=OutputState,   # 输出过滤:只返回 final_content
)

# add_node:添加节点(直接传函数,节点名自动取函数名)
builder.add_node(generate_outline)     # 节点1:生成大纲
builder.add_node(generate_draft)       # 节点2:生成初稿
builder.add_node(polish_content)       # 节点3:润色稿件
builder.add_node(finalize_content)     # 节点4:生成最终稿

# add_edge:添加边(连接节点,直线流程)
builder.add_edge(START, "generate_outline")              # 开始 → 生成大纲
builder.add_edge("generate_outline", "generate_draft")   # 大纲 → 生成初稿
builder.add_edge("generate_draft", "polish_content")     # 初稿 → 润色
builder.add_edge("polish_content", "finalize_content")   # 润色 → 最终稿
builder.add_edge("finalize_content", END)                # 最终稿 → 结束

# compile:编译工作流
chain = builder.compile()


# ---------- 9. 使用工作流 ----------
result = chain.invoke({"topic": "人工智能的未来发展"})
print("=" * 50)
print("最终创作结果:")
print("=" * 50)
print(result["final_content"])
print("=" * 50)

运行结果(示意,实际内容随模型与参数变化):

**************************************************
内容大纲生成中...

大纲已生成:
一、技术突破:从大模型到通用智能
二、社会影响:效率革命与伦理挑战
**************************************************
生成初稿中...

初稿已生成:
(两个标题下各三句话的初稿内容)
**************************************************
文章润色中...

润色完成,内容如下:
(更精炼的两个小节内容)
==================================================
最终创作结果:
==================================================
人工智能的未来发展

一、技术突破:从大模型到通用智能
……

二、社会影响:效率革命与伦理挑战
……
==================================================

💡 这条链四步串行,每一步只干一件事,这就是单一职责带来的好处:哪一步效果差,就单独改那一步的提示词,不用动整张图。注意 result 里只有 final_content——中间的 outline、draft、polished_draft 全被 output_schema 过滤掉了。

3.2 并行化模式(Parallelization)

3.2.1 概念

并行化是指多个任务同时进行,提高效率,最终汇总结果。在多角度处理同一问题时,常用这个模式。

例如要研发一款主打城市通勤的智能电动自行车,具有导航、社交、防盗等功能。在开始研发前,需要进行多维度分析:

维度分析要点
市场分析用户关注续航里程、车身重量、防盗能力,并对"骑行社交"(组队、分享路线)有新兴兴趣
竞品分析传统品牌车型智能化不足;互联网品牌车型续航和线下售后服务是短板
技术分析评估更轻量化的电池材料与车身设计以提升续航和便携性,并开发基于 GPS 和移动网络的智能防盗系统与社交功能 App 的集成

最后汇总分析结果。而并行分析不仅省时,还能提升决策质量。

📷 插图位置:并行化模式示意图(START 同时扇出到 市场分析 / 竞品分析 / 技术分析,再扇入到 汇总报告 → END)

3.2.2 模式实践

代码示例:

from typing import TypedDict
from langgraph.constants import START, END
from langgraph.graph import StateGraph


# ---------- 1. 定义状态 ----------
class AnalysisState(TypedDict):
    concept: str        # 概念
    market: str         # 市场分析
    competitor: str     # 竞品分析
    tech: str           # 技术分析
    report: str         # 汇总报告


# ---------- 2. 定义三个【并行】分析任务 ----------
def market_task(state: AnalysisState):
    """市场分析"""
    return {"market": "用户关注续航、重量、防盗,对骑行社交有兴趣..."}


def competitor_task(state: AnalysisState):
    """竞品分析"""
    return {"competitor": "传统品牌智能化不足,互联网品牌续航和售后差..."}


def tech_task(state: AnalysisState):
    """技术分析"""
    return {"tech": "轻量化电池车身、GPS防盗、社交App集成..."}


# ---------- 3. 汇总结果 ----------
def combine_results(state: AnalysisState):
    """生成最终报告"""
    report = "产品分析报告\n\n"
    report += f"市场分析:\n{state['market']}\n\n"
    report += f"竞品分析:\n{state['competitor']}\n\n"
    report += f"技术分析:\n{state['tech']}\n\n"
    report += "建议:聚焦续航、防盗、社交功能的平衡发展"
    return {"report": report}


# ---------- 4. 构建工作流 ----------
builder = StateGraph(AnalysisState)

# add_node:添加节点
builder.add_node("market", market_task)
builder.add_node("competitor", competitor_task)
builder.add_node("tech", tech_task)
builder.add_node("combine", combine_results)

# add_edge:添加边
# 并行执行三个分析:START 同时指向三个节点(扇出 fan-out)
builder.add_edge(START, "market")
builder.add_edge(START, "competitor")
builder.add_edge(START, "tech")

# 三个并行任务都完成后,汇合到 combine(扇入 fan-in)
builder.add_edge("market", "combine")
builder.add_edge("competitor", "combine")
builder.add_edge("tech", "combine")

builder.add_edge("combine", END)

# compile:编译工作流
workflow = builder.compile()


# ---------- 5. 使用 ----------
result = workflow.invoke({"concept": "城市通勤智能电动自行车"})
print(result["report"])

运行结果:

产品分析报告

市场分析:
用户关注续航、重量、防盗,对骑行社交有兴趣...

竞品分析:
传统品牌智能化不足,互联网品牌续航和售后差...

技术分析:
轻量化电池车身、GPS防盗、社交App集成...

建议:聚焦续航、防盗、社交功能的平衡发展

并行化靠"边的形状"实现,两个关键动作:

动作做法效果
扇出(fan-out)START 同时 add_edge 到 3 个节点三个分析同时开跑
扇入(fan-in)3 个节点都 add_edge 到 combine全部完成后才执行汇总

💡 扇入节点会自动等待所有上游分支完成再执行,所以 combine_results 里能放心地同时读 market、competitor、tech 三个字段,不用自己写"等三个都好了"的判断。另外这三个分支各自写的是不同字段,所以即便并行也不会互相覆盖。

3.3 路由模式(Routing)

3.3.1 概念

路由模式也被称为"智能分流":根据输入内容决定执行哪个分支。最典型的案例就是智能客服系统,可以按用户问题自动分类处理。

3.3.2 模式实践

实现一个智能客服系统,根据用户问题自动分类,达到精准匹配处理能力。核心设计在条件路由机制:

设计点说明
动态路径选择基于 LLM 分析结果动态决定执行路径(结构化返回)
分支隔离不同类型的问题由专用处理器处理

代码示例:

from langchain.chat_models import init_chat_model
from langgraph.constants import START, END
from langgraph.graph import StateGraph
from typing_extensions import Literal, TypedDict
from pydantic import BaseModel, Field


# ---------- 1. 定义状态 ----------
class State(TypedDict):
    input: str          # 用户输入
    decision: str       # 路由决策
    output: str         # 最终输出


# ---------- 2. 定义路由决策的数据结构(结构化输出)----------
class Route(BaseModel):
    # Literal:限定取值范围,模型只能在这三个里选
    step: Literal["pre_sale", "after_sale", "technical"] = Field(
        description="根据用户问题类型决定路由到售前、售后还是技术处理"
    )


# ---------- 3. 路由决策节点 ----------
def model_call_router(state: State):
    """分析用户输入,决定问题类型"""
    model = init_chat_model(model="deepseek-chat", model_provider="deepseek")
    # with_structured_output:让模型按 Route 结构返回(只输出 pre_sale/after_sale/technical)
    decision = model.with_structured_output(Route).invoke(state["input"])
    return {"decision": decision.step}


# ---------- 4. 三个不同的处理节点 ----------
def pre_sale_handler(state: State):
    """处理售前咨询"""
    return {"output": "售前咨询已处理,处理内容....."}


def after_sale_handler(state: State):
    """处理售后问题"""
    return {"output": "售后问题已处理,处理内容....."}


def technical_handler(state: State):
    """处理技术问题"""
    return {"output": "技术问题已处理,处理内容....."}


# ---------- 5. 路由函数(根据决策返回下一个节点名)----------
def route_decision(state: State):
    if state["decision"] == "pre_sale":
        return "pre_sale_handler"       # 去售前处理节点
    elif state["decision"] == "after_sale":
        return "after_sale_handler"     # 去售后处理节点
    elif state["decision"] == "technical":
        return "technical_handler"      # 去技术处理节点


# ---------- 6. 构建路由工作流 ----------
router_builder = StateGraph(State)

# add_node:添加处理节点(统一写成"节点名, 执行对象")
router_builder.add_node("pre_sale_handler", pre_sale_handler)
router_builder.add_node("after_sale_handler", after_sale_handler)
router_builder.add_node("technical_handler", technical_handler)
router_builder.add_node("model_call_router", model_call_router)

# add_edge:先经过路由决策
router_builder.add_edge(START, "model_call_router")

# add_conditional_edges:条件边 —— 根据路由函数的结果选择走哪个分支
router_builder.add_conditional_edges(
    "model_call_router",                                        # 源节点
    route_decision,                                             # 路由函数(返回下一个节点名)
    ["pre_sale_handler", "after_sale_handler", "technical_handler"],  # 可选分支
)

# add_edge:所有分支最终都结束
router_builder.add_edge("pre_sale_handler", END)
router_builder.add_edge("after_sale_handler", END)
router_builder.add_edge("technical_handler", END)

# compile:编译工作流
router_workflow = router_builder.compile()


# ---------- 7. 测试 ----------
test_cases = [
    "我想了解一下你们产品的价格和功能",                    # 售前咨询
    "我购买的产品有质量问题,需要退货",                     # 售后问题
    "这个软件安装后无法正常运行,报错代码0x80070005",        # 技术问题
    "请问你们的售后服务政策是什么",                        # 售后问题
    "我的订单已经发货但还没收到",                          # 售后问题
    "如何配置数据库连接参数",                              # 技术问题
]

for test_case in test_cases:
    print("*" * 50)
    result = router_workflow.invoke({"input": test_case})
    print(f"用户问题:{test_case}\n{result['output']}")

运行结果:

**************************************************
用户问题:我想了解一下你们产品的价格和功能
售前咨询已处理,处理内容.....
**************************************************
用户问题:我购买的产品有质量问题,需要退货
售后问题已处理,处理内容.....
**************************************************
用户问题:这个软件安装后无法正常运行,报错代码0x80070005
技术问题已处理,处理内容.....
**************************************************
用户问题:请问你们的售后服务政策是什么
售后问题已处理,处理内容.....
**************************************************
用户问题:我的订单已经发货但还没收到
售后问题已处理,处理内容.....
**************************************************
用户问题:如何配置数据库连接参数
技术问题已处理,处理内容.....

💡 路由模式的关键在Route 这个 Pydantic 模型:Literal["pre_sale", "after_sale", "technical"] 把模型的输出锁死在三个取值里,所以 route_decision 里的三个分支不会漏——不会出现"模型返回了一个没准备的类型"导致路由失败。这也解释了为什么结构化输出和条件边是天生一对。

3.4 协调者-工作者模式(Orchestrator-Workers)

3.4.1 概念

协调者-工作者模式可以理解为一个大脑(协调者)分配任务,多个工人(工作者)执行,最后合成最终结果。例如处理几百页的技术文档时,可以让协调者拆分文档,接着安排多个工作者并行处理不同章节,最终汇总结果。

3.4.2 模式实践

三个角色分工如下:

角色职责
协调者根据 {topic} 生成报告大纲,并把生成内容的子任务指派给工作者
工作者生成大纲对应的内容:协调者生成 3 个标题就需要 3 个工作者,生成 10 个标题就需要 10 个工作者
合成器汇总所有工作者的成果

关键在于任务指派:协调者在运行时才需动态分配工作者,也就是说边的数量在运行时才能确定。为了支持这种模式,LangGraph 支持从条件边返回 Send 对象:

Send 的参数含义
第 1 个节点的名称(要把任务发给哪个节点)
第 2 个要传递给该节点的状态

代码示例:

from langchain.chat_models import init_chat_model
from langgraph.constants import START, END
from langgraph.graph import StateGraph
from langgraph.types import Send
from typing import Annotated, TypedDict, List
import operator
from pydantic import BaseModel


# ---------- 1. 定义状态 ----------
class State(TypedDict):
    topic: str
    sections: list                                     # 协调者生成的计划
    # Annotated[list, operator.add]:reducer —— 多个工作者的结果自动【合并/追加】
    completed_sections: Annotated[list, operator.add]  # 工作者完成的结果
    final_report: str


# ---------- 2. 定义数据结构(结构化输出)----------
class Section(BaseModel):
    name: str
    description: str


class Sections(BaseModel):
    sections: List[Section]


# ---------- 3. 创建规划器(让模型按 Sections 结构输出)----------
model = init_chat_model(model="deepseek-chat", model_provider="deepseek")
planner = model.with_structured_output(Sections)


# ---------- 4. 工作者私有状态(由 Send 传入,不在 State 中)----------
class WorkerState(TypedDict):
    section: Section      # 单个工作者领到的任务


# ---------- 5. 协调者节点 —— 制定计划 ----------
def orchestrator(state: State):
    """协调者:分析任务并制定执行计划"""
    report_sections = planner.invoke(
        f"为主题'{state['topic']}'制定报告大纲,包含3个章节"
    )
    return {"sections": report_sections.sections}


# ---------- 6. 工作者节点 —— 执行具体任务 ----------
def llm_call(state: WorkerState):
    """工作者:根据分配的任务生成内容"""
    section = state["section"]      # 从协调者接收的任务
    result = model.invoke(
        f"编写报告章节:{section.name},内容要求:{section.description}"
    )
    return {"completed_sections": [result.content]}   # 结果会自动合并


# ---------- 7. 汇总节点 ----------
def synthesizer(state: State):
    """汇总所有工作者的成果"""
    completed_sections = state["completed_sections"]
    final_report = "\n\n---\n\n".join(completed_sections)
    return {"final_report": final_report}


# ---------- 8. 任务分配函数 —— 关键部分!----------
def assign_workers(state: State):
    """为每个任务创建工作作者(动态生成多个并行任务)"""
    worker_tasks = []
    for section in state["sections"]:
        # Send("节点名", {传给该节点的状态}):把任务【动态发送】给工作者节点
        worker_tasks.append(Send("llm_call", {"section": section}))
    return worker_tasks


# ---------- 9. 构建工作流 ----------
builder = StateGraph(State)
builder.add_node("orchestrator", orchestrator)
builder.add_node("llm_call", llm_call)
builder.add_node("synthesizer", synthesizer)

builder.add_edge(START, "orchestrator")

# 关键:协调者后创建多个工作者(条件边 + Send 实现动态并行分发)
builder.add_conditional_edges(
    "orchestrator",
    assign_workers,
    ["llm_call"],       # 创建的工作者都指向 llm_call 节点
)

# 所有工作者完成后汇总
builder.add_edge("llm_call", "synthesizer")
builder.add_edge("synthesizer", END)

worker = builder.compile()

response = worker.invoke({"topic": "中国近代史"})
print(response)

运行结果:

{'topic': '中国近代史', 'sections': [Section(name='第一章:晚清变局与救亡图存的探索(1840—1911)', description='……'), Section(name='第二章:共和初建与革命新道路的探索(1912—1949)', description='……'), Section(name='第三章:……', description='……')], 'completed_sections': ['(第一章正文)', '(第二章正文)', '(第三章正文)'], 'final_report': '(第一章正文)\n\n---\n\n(第二章正文)\n\n---\n\n(第三章正文)'}

💡 这里是"边在运行时才确定"的最直观例子:assign_workers 遍历 sections 列表,有几个章节就返回几个 Send——3 个章节就是 3 条并行分支,10 个章节就是 10 条。所以 llm_call 读到的 state["section"] 并不是 State 里的字段,而是 Send 第二参数单独传进去的私有状态(这正是 2.3 节讲的那套机制:能传,但不会出现在最终状态里)。另外 completed_sections 用 operator.add 做 reducer,多个工作者的结果才会自动合并成列表,而不是互相覆盖。


四、复盘(附答案)

💡 思考题

  1. Overwrite 解决的是什么问题?不加它,replace_messages 节点的返回会怎样被处理?

  2. input_schema、output_schema、状态类本身(state_schema)分别管什么?为什么内部状态可以比对外接口更复杂?

  3. 私有状态是怎么做到"节点之间能传、最终输出里没有"的?Node1Output 和 Node2Input 能不能合并成一个类?

  4. 并行化模式里,扇出和扇入分别靠什么实现?扇入节点需要自己判断"上游都跑完了"吗?

  5. 路由模式里,Route 模型为什么要用 Literal 限定取值范围?

  6. Send 的第一个参数和第二个参数分别是什么?为什么说协调者-工作者模式"边的数量在运行时才确定"?

📝 答案

  1. Overwrite 解决的是"某些场景下必须整体覆盖、而不是合并"的问题,典型就是清空聊天记录重新开始。messages 字段挂着 operator.add 这个 reducer,节点的返回值会被当作"要追加的内容"合并进旧列表;不加 Overwrite 时结果是 ['initial', 'first message', 'replacement message'],加上之后变成 ['replacement message']——强制绕过 reducer,直接整体替换。

  2. input_schema 验证调用图形时的输入结构(只接受声明的字段);output_schema 过滤图的输出(只返回声明的字段);状态类本身是图内部流转的完整状态,节点之间通信用。所以内部状态可以比对外接口更复杂——它可以装着中间结果、上下文、临时数据,而对外只暴露该暴露的字段,既保护内部实现,也让接口稳定。

  3. 靠节点的类型注解:node_1 的返回注解是 Node1Output、node_2 的参数注解是 Node2Input,两者有同名同类型字段 sensitive_data,于是 node_1 交出的数据正好被 node_2 读走;而这个字段不在 OverallState 里,所以不会出现在最终输出中(最终只剩 final_result)。能合并成一个类——字段一样,写成同一个 SensitiveData 完全等价;拆成两个类只是让"谁的输入 / 谁的输出"语义更清晰。

  4. 扇出靠从同一个起点加多条边:START 同时连到 market、competitor、tech,三个节点同时开跑;扇入靠多个节点连到同一个节点:三个分析节点都连到 combine。不需要自己判断——扇入节点会自动等待所有上游分支完成再执行,所以 combine_results 里可以直接读三个字段。

  5. 因为结构化输出需要约束取值范围。Literal["pre_sale", "after_sale", "technical"] 让模型只能在这三个值里选,于是 route_decision 里的三个分支覆盖了全部可能,不会出现"模型返回了第四个值、路由函数却没有对应分支"导致流程中断的情况。条件边 + 结构化输出因此是天然搭配。

  6. 第一个参数是节点的名称(任务发给哪个节点),第二个参数是要传递给该节点的状态(这里是 {"section": section})。之所以说"边的数量在运行时才确定",是因为 assign_workers 是在运行时遍历 state["sections"] 才能知道有几个章节,然后返回几个 Send——3 个章节就 3 条并行分支,10 个章节就 10 条,图结构里写死的只有"协调者 →(条件边)→ llm_call"这一条。


🎯 闭幕

​

如果本文对你有帮助,欢迎:

👍 点赞 | ⭐ 收藏 | 👤 关注作者 | 💬 留言交流你的疑问或补充

你的每一次互动都是我继续更新的动力,我们下一篇见!🚀

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/2601_96587588/article/details/167127657

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--