LangGraph 工作流构建(全)
目录
一、LangGraph 的核心组件
State(状态):[数据传输],所有节点的“共享工作记忆”,存储任务中的全部信息。
特点:增量更新(节点只返回部分修正)
Node(节点):[最小执行单元],例如调用AI、查询数据库等。
节点职责通常单一,输入当前 State,并返回 State 的部分更新
Edge(边):[定义工作流执行顺序],普通边(固定顺序执行)、条件边(条件判断)
支持:循环、回溯、并行、多分支、终止(END)
StateGraph(状态图):[工作流的蓝图] 定义了状态的结构、节点和边的整体布局
所有节点和边都必须在一个 StateGraph 实例上定义和组合
二、核心组件应用示例
1. State状态(定义参数)
State需要那些字段完全基于业务自定义。
messages: Annotated[list, operator.add] 是常用字段类型,可以累加所有节点需要记录的内容。
messages 存储那些数据,一般有2种场景:
场景一、LLM对话类一般存储,每个节点和大模型的对话记录:
{"role": "user", "content": "用户输入"} {"role": "assistant", "content": "AI回答"} # 例如返回当前节点结果【只返回需要更新的字段(字典),LangGraph会自动合并到状态中】 return { "messages": [ {"role": "assistant", "content": “大模型返回的结果文本”} ], "next_step": "bb" }场景二、执行日志 / 流程记录(每个节点记一句干了什么、结果是什么,方便排查):
def aa_node(state: ConstructionState): return { "messages": [ { "step": "节点名称", "action": "节点作用", "status": "成功/失败", "data": "上一节点参数值,当前节点处理结果" } ], "next_step": "bb" }
class ConstructionState(TypedDict):
# 消息上下文:自动累加追加,不覆盖历史
# 作用:所有节点都能看到前面所有步骤的历史
messages: Annotated[list, operator.add] # Annotated[list, operator.add] 创建一个累加的list参数
# 【路由控制】控制条件边下一步走向
# 可选值:"bb", "cc", "dd", "end" 等任意节点名称
# 判断函数根据该字段决定流程走向
next_step: str
# 用户原始输入内容
input_text: str
# 工作流最终输出结果(当前节点输出的结果)
result: str
# ...
2. Node普通节点
定义三个普通函数方法节点
# ========================
# 节点函数 伪代码
# 规则:函数必须 return 字典格式 ConstructionState → 不返回的字段节点参数值不变。
# ========================
def aa_node(state: ConstructionState):
"""
第一个节点:入口节点
作用:处理输入、设置路由、记录消息
"""
# 1. 读取用户输入(所有state字段都能拿到)
input_text = state["input_text"]
# 2. 你的业务逻辑(随便写)
new_msg = f"[aa节点] 处理用户输入,{input_text} 将问题控制在20字内。"
result_llm = "调用大模型,返回的接口"
# 3. 返回要更新的state字段(只返回你要改的!)
return {
"messages": [new_msg], # 自动累加
"next_step": "bb", # 控制下一步走向 bb
"result": result_llm # 节点结果
}
def bb_node(state: ConstructionState):
"""
第二个节点:中间处理节点
作用:业务处理、加工数据
"""
# 1. 读取上一步状态
input_text = state["input_text"]
prev_result = state["result"]
# 2. 业务逻辑
print("=== 执行 bb 节点 ===")
new_msg = f"[bb节点] 基于aa结果继续处理:{prev_result}"
new_result = f"{prev_result} → bb节点处理完成"
# 3. 返回更新
return {
"messages": [new_msg],
"result": new_result
}
def cc_node(state: ConstructionState):
# 和以上写法一样
2.1. Node工具调用节点
from langgraph.prebuilt import ToolNode
from langchain_core.tools import tool
# 1. 定义工具
@tool
def get_weather(city: str) -> str:
"""获取天气"""
return f"{city}天气晴朗"
@tool
def calculate(expression: str) -> str:
"""计算数学表达式"""
return str(eval(expression))
tools = [get_weather, calculate]
# 2. ToolNode 是 LangGraph 预置的工具节点(自动执行工具调用,返回最后一条结果)
# 注意:ToolNode 可以调用多个工具,但是不会做循环调用(react)
tool_node = ToolNode(tools)
# 2. 在图中使用(和普通节点使用方式一样)
graph_builder.add_node("tools", tool_node) # 直接添加,不需要写函数
3. Edge边的两种写法
StateGraph状态图(state状态)、添加node节点、Edge边的完整使用案例。
3.1. 普通边(固定顺序)
开始 → aa → bb → cc → 结束
def create_construction_workflow() -> StateGraph:
# 1.定义状态图(定义工作流)【类似找白纸,后续在上边画节点】
workflow = StateGraph(ConstructionState) # 参数(State 状态)
# 2.添加节点
workflow.add_node("aa", aa_node) # (定义节点命名,对应的函数方法)
workflow.add_node("bb", bb_node)
workflow.add_node("cc", cc_node)
# ========================
# 3.普通边(连接节点),执行顺序aa节点 → bb节点 → cc节点
# ========================
workflow.set_entry_point("aa") # 设置入口节点为aa
workflow.add_edge("aa", "bb") # aa执行完执行bb节点
workflow.add_edge("cc", END) # 设置结束节点为bb
return workflow.compile() # 编译工作流
3.2. 条件边(条件判断)
开始 → aa → 分支判断
├─ go_bb → bb → cc → 结束
└─ go_cc → cc → 结束
# 判断函数:aa 执行完,决定去 bb 还是 cc
def judge_aa_to_where(state):
if state["next_step"] == "go_bb":
return "go_bb" # 对应去 bb
else:
return "go_cc" # 对应去 cc
def create_construction_workflow() -> StateGraph:
# 1. 定义工作流(状态)
workflow = StateGraph(ConstructionState)
# 2. 添加节点
workflow.add_node("aa", aa_node) # (定义节点命名,对应的函数方法)
workflow.add_node("bb", bb_node)
workflow.add_node("cc", cc_node)
# 3. 设置边入口
workflow.set_entry_point("aa")
# ========================
# 条件边(多分支判断):aa 执行完 → 判断走 bb 还是 cc节点
# 注意:一次 add_conditional_edges 只能设置【一个节点】的分支跳转。
# ========================
workflow.add_conditional_edges(
"aa", # 从 aa节点 出发
judge_aa_to_where, # 判断函数
{ # 路线表
"go_bb": "bb", # 如果judge_aa_to_where函数返回go_bb,就执行bb节点
"go_cc": "cc" # 如果judge_aa_to_where函数返回go_cc,就执行cc节点
}
)
# bb 走完 → 去 cc
workflow.add_edge("bb", "cc")
# cc 走完 → 结束
workflow.add_edge("cc", END)
return workflow.compile() # 编译工作流
4. 执行工作流
app = create_construction_workflow();# 获取工作流实例
result = app.invoke(initial_state) # 执行工作流
三、LangGraph异常处理(重试、熔断)
完整代码路径(超时控制、重试机制、异步调用):\week09\4\p34常见异步陷阱及规避
1. 重试
1.1. 节点重试
from langgraph.types import RetryPolicy
# 基本重试策略
basic_policy = RetryPolicy(
max_attempts=3, # 最大重试次数(1+3=4次)
retry_delay=1.0, # 重试延迟(秒)
backoff_factor=2.0, # 退避因子,每次等待时间翻倍(2→4→8秒)
max_delay=60.0, # 最大延迟时间
retry_on=(ValueError,)# 可重试的异常类型
)
...
# retry=basic_policy 加入节点重试规则
builder.add_node("get_user_node", get_user_node, retry=basic_policy)
...
1.2. 工具节点重试
from langgraph.types import RetryPolicy
from langgraph.prebuilt import ToolNode
# 工具节点
tool_node = ToolNode(
# 配置工具列表
tools=[api_call_function],
# 工具执行失败的重试策略(仅对工具本身生效)
retry_policy=RetryPolicy(
max_attempts=3, # 最多重试3次(总共执行1+3=4次)
retry_delay=2.0, # 第一次重试前等待2秒
backoff_factor=2.0, # 退避因子=2,每次等待时间翻倍(2→4→8秒)
retry_on=(ConnectionError, TimeoutError) # 只在这两种异常时重试
)
)
# 将工具节点添加到流程图中
builder.add_node("api_call", tool_node)
1.3. 熔断机制
# 熔断器三种状态 正常 → 失败重试次数超阈值 → 熔断开启(直接拒绝请求)→ 半开尝试 → 恢复/继续熔断s四
四、带人工确认节点的工作流(工作流持久化)
人工确认节点依赖工作流持久化:将每次节点执行后的 State 存入数据库(即工作流持久化),人工确认节点才能中断后恢复执行。
核心用途:
- 中断恢复 — 在 自定义节点前暂停(例如在confirm节点前暂停),等人确认后继续执行(持久化的核心用途)
- 进度查询 — 查看当前执行到哪个节点(不持久化也能做,但服务重启后查不到)
- 历史追溯 — 查看某次请求各节点的执行记录,用于排查问题
D:\dev\python-projects\demo-01-project\project-fs-demo\graph\project_schedule\workflow.py
工作流添加持久化:
# 配置工作流状态存储
# 这里使用 SQLite 数据库存储状态,也可以选择其他存储方式(如 Redis、MongoDB 等)
memory = SqliteSaver(
sqlite3.connect("checkpoints.db", check_same_thread=False)
)
# 编译工作流(到达 confirm 节点前自动 return,等待下次 invoke 时从中断处继续)
return workflow.compile(interrupt_before=["confirm"], checkpointer=memory)
D:\dev\python-projects\demo-01-project\project-fs-demo\app.py
# 首次请求接口:提交用户问题,启动工作流
# 需要人工确认 → 返回 need_confirm(含 thread_id);不需要 → 返回 success
# 首次请求接口:提交用户问题,启动工作流
# 需要人工确认 → 返回 need_confirm(含 thread_id);不需要 → 返回 success
@app.post("/api/add/progress", response_model=QueryResponse)
async def handle_add_progress(request: QueryRequest):
if not request.user_query or not request.user_query.strip():
raise HTTPException(status_code=400, detail="查询参数不能为空")
try:
# 初始化State状态值
initial_state = AddProgressState(
user_query=request.user_query,
original_form_info=None,
uploaded_file=None,
messages=[]
)
"""
生成唯一的工作流ID,用于中断恢复(人工确认后继续执行)
"""
config = {"configurable": {"thread_id": str(uuid.uuid4())}}
"""
执行工作流
"""
result = workflow.invoke(initial_state.model_dump(), config=config)
"""
检查工作流返回结果是否需要人工确认
"""
state = workflow.get_state(config)
if state.next == ("confirm",):
return QueryResponse(
status="need_confirm",
result="请确认以下信息是否正确",
messages=result["confirm_data"],
thread_id=config["configurable"]["thread_id"]
)
return QueryResponse(
status="success",
result=result.get("result", "操作完成"),
)
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
# 二次确认接口:用户确认/拒绝后,从中断处恢复执行
# 用户确认(yes)→ 继续推到飞书;用户拒绝(no)→ 工作流正常结束
# 二次确认接口:用户确认/拒绝后,从中断处恢复执行
# 用户确认(yes)→ 继续推到飞书;用户拒绝(no)→ 工作流正常结束
@app.post("/api/add/progress/confirm", response_model=QueryResponse)
async def handle_confirm(request: ConfirmRequest):
try:
config = {"configurable": {"thread_id": request.thread_id}}
workflow.update_state(config, {"user_confirmed": request.user_confirmed})
result = workflow.invoke(None, config=config)
# 执行完成后可以删除工作流状态
#await workflow.checkpointer.delete_thread(request.thread_id)
return QueryResponse(
status="success",
result=result["result"],
)
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
五、构建“反思机制”智能体
反思智能体 = 在图里多加一个「反思/审稿」节点:通过就往下走,不通过就带着“批评问题”打回「生成」节点重新生成问题。
它解决什么问题?
- 普通 Agent:问一次 → 答一次,容易漏点、格式乱、逻辑跳步。
- 反思模式:生成 → 批评 → 改 → 再批评,用多轮自评提高质量(写作、代码、复杂推理等场景)。
和 ReAct 有啥区别?
ReAct = Reasoning(推理)+ Acting(行动)
方式:[思考] → 调工具 → [看结果] → [思考] → 调工具 → … → 最终回答
反思智能体 = 生成答案 → 批评自己的答案 → 不满意就重写
方式:任务 → [生成草稿/回答] → [反思:哪里不好] → 通过?结束 : 带着批评再生成
流程长什么样?

- 生成节点:根据任务 + 上一轮批评,产出/修改草稿(回答)
- 反思节点:只负责挑毛病、给建议,不直接替用户定稿
- 条件边:通过继续;未通过打回,并限制最大轮数防死循环
State 里通常放什么?
user_request:原始任务draft:当前草稿(节点输出的回答)critique:反思意见(打回时给生成节点用)iteration:已迭代次数(配合上限)
什么时候用 / 不用?
- 适合:质量比速度重要;有明确评判标准(格式、单测、rubric)
- 不太适合:简单问答、强实时对话(成本高、延迟大);纯检索更适合 RAG + 工具
更多推荐



所有评论(0)