在Python中,这种复杂的workflow编排可以使用 LangGraph 框架来实现。LangGraph 是 LangChain 团队推出的工作流编排框架,核心思想是将业务流程建模为一张有向图:
- 节点(Node):每个节点是一个处理步骤(一个函数)
- 边(Edge):定义节点之间的执行顺序
- 状态(State):在节点之间传递的共享数据
LangGraph 之于 Python,就像 Spring AI Alibaba 之于 Java。两者的核心概念对应关系:
要是用Langgraph必须首先安装依赖:
``` pip install langgraph ```
接下来分步学习LangGraph的基本使用。
定义状态
状态是所有节点共享的数据容器,使用 TypedDict 定义,TypedDict就等价于Spring AI Alibaba中的OverAllState:
```
from typing import TypedDict, NotRequired
class MyState(TypedDict):
"""工作流状态,所有节点通过读写状态来传递数据"""
user_input: str # 用户输入
result: str # 处理结果
score: int # 评分
```
注意:这里和普通的Python类不同!
学过Python看到 class MyState: 里面写 user_input: str 可能会条件反射地认为这是类变量(类似Java里的 static 字段,所有实例共享同一个值)。这是错觉。
TypedDict 是Python专门用来”描述字典结构”的语法糖。它里面写的所有 字段名: 类型 都不是真正的变量赋值,只是给字典字段做的类型声明——告诉编辑器和类型检查器”这种字典里应该有这些字段,每个字段是什么类型”。
举两个对比就清楚了:
```
# ❌ 这才是真正的Python类,里面是"类变量"
class NormalClass:
user_input: str = "默认值" # 有 = 赋值,这是类属性
# ✅ TypedDict 里没有 = 赋值,只是"字段声明"
class MyState(TypedDict):
user_input: str # 没有 =,只是说"这种字典里有个叫user_input的key,值是str类型"
```
用法上 MyState 就是一个字典:
```
# 创建一个MyState"实例"——其实就是创建一个普通字典
state: MyState = {"user_input": "你好"}
# 读取:用字典的中括号语法,不是 . 语法
print(state["user_input"]) # ✅ 正确
# print(state.user_input) # ❌ 错误!TypedDict不是类,没有属性访问
# 写入:和字典完全一样
state["result"] = "已处理"
```
所以在节点函数里你会看到的都是 state["xxx"] 这种字典访问写法,不是 state.xxx。
定义节点
节点就是一个普通的Python函数,接收当前状态,返回要更新的字段:
```
def process_node(state: MyState) -> dict:
"""节点函数:读取状态中的数据,处理后返回要更新的字段"""
user_input = state["user_input"] # 从状态中读取
result = f"已处理:{user_input}"
return {"result": result, "score": 85} # 返回要写入状态的字段
```
节点函数的规则:
- 参数:接收整个状态字典
- 返回值:一个字典,包含要更新的状态字段(不需要返回所有字段,只返回有变化的)
定义图
使用 StateGraph 创建图,添加节点:
```
from langgraph.graph import StateGraph, START, END
# 创建图,指定状态类型
graph_builder = StateGraph(MyState)
# 添加节点:(节点名称, 节点函数)
graph_builder.add_node("process", process_node)
graph_builder.add_node("evaluate", evaluate_node)
```
添加边
无条件边: 节点A执行完后,一定执行节点B
```
# START → process → evaluate → END
graph_builder.add_edge(START, "process")
graph_builder.add_edge("process", "evaluate")
graph_builder.add_edge("evaluate", END)
```
条件边: 根据状态中的数据决定下一步走哪个节点
```
def route_function(state: MyState) -> str:
"""路由函数:根据状态返回下一个节点的名称"""
if state.get("score", 0) >= 80:
return "pass"
return "retry"
# 从evaluate节点出发,根据route_function的返回值决定走哪条边
graph_builder.add_conditional_edges(
"evaluate", # 源节点
route_function, # 路由函数
{
"pass": "finish", # 路由函数返回"pass" → 走finish节点
"retry": "process", # 路由函数返回"retry" → 走process节点
}
)
```
并行边: 一个节点执行完后,同时启动多个节点
```
# geocode执行完后,同时启动三个节点
graph_builder.add_edge("geocode", "search_activities")
graph_builder.add_edge("geocode", "search_restaurants")
graph_builder.add_edge("geocode", "search_hotels")
# 三个并行节点都完成后,汇聚到assemble节点
graph_builder.add_edge(
["search_activities", "search_restaurants", "search_hotels"],
"assemble",
)
```
执行图
```
# 创建图,指定状态类型
graph_builder = StateGraph(MyState)
# 添加节点:(节点名称, 节点函数)
graph_builder.add_node("process", process_node)
graph_builder.add_node("evaluate", evaluate_node)
# 编译图
graph = graph_builder.compile()
# 非流式执行:一次性返回整个字典
result = graph.invoke({"user_input": "你好"})
print(result) # {"user_input": "你好", "result": "已处理:你好", "score": 85}
# 流式执行:逐个节点返回中间结果
for chunk in graph.stream({"user_input": "你好"}, stream_mode="updates"):
print(chunk)
```
graph.stream 返回的 chunk 到底是什么?
很多同学第一次看到 for chunk in graph.stream(...) 会以为每个 chunk 是一段文本(类似LLM流式输出那种逐字返回的字符串)。不是——这里的 chunk 是一个字典,它记录了刚刚执行完的那个节点的更新。
具体结构是:
```
{
"节点名": {
# 这个节点本次执行时 return 的内容
"字段1": 值1,
"字段2": 值2,
...
}
}
```
举个具体例子。假设图里有 process 和 evaluate 两个节点,依次执行:
```
for chunk in graph.stream({"user_input": "你好"}, stream_mode="updates"):
print(chunk)
# 第1次循环(process 节点执行完):
# {"process": {"result": "已处理:你好", "score": 85}}
# ↑ key是节点名 ↑ value是这个节点 return 的字典
# 第2次循环(evaluate 节点执行完):
# {"evaluate": {"score": 90}}
# ↑ key是节点名 ↑ evaluate节点只更新了score字段
```
Human-in-the-Loop
有些步骤需要暂停执行,等待用户输入后再继续。LangGraph通过 interrupt() 实现:
```
from langgraph.types import interrupt, Command
from langgraph.checkpoint.memory import MemorySaver
def ask_user_node(state: MyState) -> dict:
"""需要用户输入的节点"""
# interrupt() 会暂停图的执行,将参数作为中断值返回给调用方
user_answer = interrupt("请输入你的偏好:")
# 用户恢复执行后,interrupt() 返回用户的输入
return {"user_input": user_answer}
# 使用human-in-the-loop时,必须配置checkpointer(用于保存中断时的状态)
graph = graph_builder.compile(checkpointer=MemorySaver())
# 第一次执行:遇到interrupt()会暂停
config = {"configurable": {"thread_id": "user-123"}}
for chunk in graph.stream({"user_input": "开始"}, config=config):
print(chunk) # 执行到ask_user_node时暂停,返回中断值
# 用户输入后,恢复执行
for chunk in graph.stream(
Command(resume="我喜欢安静的地方"), # 传入用户的回答
config=config, # 必须用同一个thread_id
):
print(chunk) # 从中断处继续执行
```
需要注意的是:
MemorySaver是内存中的状态保存器,用于在中断期间保存图的执行状态。thread_id用于隔离不同用户的会话。- 在使用了HIL之后,在graph执行时需要传递config参数,因为需要通过thread——id来回复中断前的状态
综合案例:外卖订单处理
下面用一个简单的案例把以上所有语法点串起来。场景:一个外卖订单处理系统的后端接口,包含3个节点。
业务流程:
receive_order:接收订单,提取菜品和地址confirm_order(HITL):向用户确认订单信息,等待用户确认或修改dispatch_order:根据确认结果决定派单或取消

图的定义
```
from typing import TypedDict, NotRequired
from langgraph.graph import StateGraph, START, END
from langgraph.types import interrupt
from langgraph.checkpoint.memory import MemorySaver
# 1. 定义状态
class OrderState(TypedDict):
raw_message: str # 用户原始消息
dish: str # 菜品名称
address: str # 配送地址
confirmed: bool # 用户是否确认
dispatch_result: str # 派单结果
# 2. 定义节点
def receive_order(state: OrderState) -> dict:
"""解析用户消息,提取菜品和地址"""
message = state["raw_message"]
# 简单模拟解析(实际项目中会用LLM)
return {"dish": "宫保鸡丁", "address": "望京SOHO"}
def confirm_order(state: OrderState) -> dict:
"""暂停执行,等待用户确认订单"""
prompt = f"您的订单:{state['dish']},配送到:{state['address']}。确认下单吗?(yes/no)"
answer = interrupt(prompt) # 暂停,等待用户输入
return {"confirmed": answer.strip().lower() in ("yes", "是", "确认")}
def dispatch_order(state: OrderState) -> dict:
"""派单"""
return {"dispatch_result": f"{state['dish']} 已派单到 {state['address']}"}
# 3. 路由函数:根据确认结果决定下一步
def route_after_confirm(state: OrderState) -> str:
if state.get("confirmed"):
return "dispatch"
return "cancel"
# 4. 构建图
builder = StateGraph(OrderState)
builder.add_node("receive_order", receive_order)
builder.add_node("confirm_order", confirm_order)
builder.add_node("dispatch_order", dispatch_order)
builder.add_edge(START, "receive_order") # 无条件边
builder.add_edge("receive_order", "confirm_order") # 无条件边
builder.add_conditional_edges( # 条件边
"confirm_order",
route_after_confirm,
{"dispatch": "dispatch_order", "cancel": END},
)
builder.add_edge("dispatch_order", END)
# 5. 编译(带checkpointer,支持HITL中断恢复)
graph = builder.compile(checkpointer=MemorySaver())
```
Web接口实现
```
import uuid
import json
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
app = FastAPI()
class RunRequest(BaseModel):
threadId: str
input: dict | None = None # 首次执行时传入初始数据
resume: str | None = None # 恢复执行时传入用户的回答
@app.post("/api/threads")
def create_thread():
"""创建新的会话线程"""
return {"threadId": str(uuid.uuid4())}
@app.post("/api/orders/stream")
def stream_order(req: RunRequest):
"""流式执行图,通过SSE推送事件"""
config = {"configurable": {"thread_id": req.threadId}}
def event_stream():
# 判断是首次执行还是恢复执行
if req.input:
graph_input = req.input # 首次:传入初始数据
else:
graph_input = Command(resume=req.resume) # 恢复:传入用户回答
# stream_mode="updates" 每个节点完成后推送一次
for chunk in graph.stream(graph_input, config=config, stream_mode="updates"):
# 检查是否遇到中断(HITL节点暂停)
if "__interrupt" in chunk:
interrupt_data = chunk["__interrupt"][0]
yield f"data: {json.dumps({'type': 'interrupt', 'value': interrupt_data['value']})}\n\n"
return # 中断后停止,等待前端下次请求恢复
# 推送节点完成事件
for node_name in chunk:
if not node_name.startswith("__"):
yield f"data: {json.dumps({'type': 'node', 'name': node_name})}\n\n"
yield f"data: {json.dumps({'type': 'done'})}\n\n"
return StreamingResponse(event_stream(), media_type="text/event-stream")
```
前端交互流程:
- 调用
POST /api/threads→ 获得threadId - 调用
POST /api/orders/stream,body:{threadId, input: {raw_message: "一份宫保鸡丁送到望京SOHO"}} - 收到SSE事件:
{type: "node", name: "receive_order"}→ 订单已解析 - 收到SSE事件:
{type: "interrupt", value: "您的订单:宫保鸡丁...确认下单吗?"}→ 图暂停了 - 前端展示确认提示,用户点击”确认”
- 调用
POST /api/orders/stream,body:{threadId, resume: "yes"}→ 恢复执行 - 收到SSE事件:
{type: "node", name: "dispatch_order"}→ 派单完成 - 收到SSE事件:
{type: "done"}→ 图执行结束