1. LangGraph与Human-in-the-Loop工作流概述
LangGraph作为LangChain生态中的工作流编排框架,其核心价值在于实现了复杂AI任务的模块化执行。而Human-in-the-Loop(人机协同)机制则是将人类判断引入自动化流程的关键设计,特别适合需要人工复核的高风险场景。比如在酒店预订系统中,当AI尝试执行预订操作时,系统可以暂停流程并等待人工确认,避免错误预订造成的损失。
与传统工作流引擎不同,LangGraph的人机协同具备两个独特优势:首先是中断持久化能力,工作流可以在任意节点暂停并保存完整上下文,即使服务器重启也不影响;其次是细粒度控制,开发者可以精确指定哪些操作需要人工介入,其余环节仍保持自动化。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 核心机制解析
2.1 中断控制原语
interrupt()函数是人机协同的核心原语,其工作原理类似于操作系统的系统调用。当工作流执行到该函数时,会触发以下动作:
- 将当前状态(包括内存、变量、调用栈)序列化到检查点
- 向外部系统发送中断请求
- 进入等待状态,释放计算资源
检查点存储支持多种后端,开发中常用InMemorySaver,生产环境建议使用PostgresSaver:
python复制from langgraph.checkpoint.postgres import PostgresSaver
from sqlalchemy import create_engine
engine = create_engine("postgresql://user:pass@localhost/db")
checkpointer = PostgresSaver(engine)
2.2 恢复机制设计
工作流恢复通过Command(resume=...)实现,支持三种恢复模式:
- accept:直接继续执行
- edit:修改输入参数后继续
- response:返回自定义响应
典型的生产级实现会构建Web界面供人工操作,后端处理逻辑如下:
python复制@app.post("/resume/{thread_id}")
async def resume_workflow(thread_id: str, action: ResumeAction):
command = Command(
resume=[{
"type": action.type,
"args": action.args
}]
)
return StreamingResponse(
agent.stream(command, {"configurable": {"thread_id": thread_id}})
)
3. 实战:电商审核系统搭建
3.1 场景设计
假设我们需要构建一个智能客服系统,其中优惠券发放需要人工审核。工作流包含:
- 用户请求解析
- 资格自动校验
- 人工审核节点
- 优惠券发放
3.2 关键实现
定义审核工具时需注意:
- 工具描述要明确说明需要人工介入
- 参数校验要严格,减少无效中断
python复制from pydantic import BaseModel
class CouponInput(BaseModel):
user_id: str
coupon_type: str
amount: float
@tool(args_schema=CouponInput)
def issue_coupon(user_id: str, coupon_type: str, amount: float):
"""发放优惠券(需要人工审核)"""
interrupt(
f"用户{user_id}申请{coupon_type}优惠券{amount}元",
config={
"approval_threshold": amount > 1000 # 大额优惠强制人工审核
}
)
return f"已发放{coupon_type}优惠券"
3.3 工作流编排
使用StateGraph构建带有人工审核节点的流程图:
python复制from langgraph.graph import StateGraph
workflow = StateGraph(AgentState)
# 添加自动节点
workflow.add_node("parse_request", parse_user_request)
workflow.add_node("check_eligibility", check_qualification)
# 添加人工审核节点
workflow.add_node("human_approval", issue_coupon)
# 设置边条件
workflow.add_conditional_edges(
"check_eligibility",
lambda x: "human_approval" if x["needs_approval"] else "auto_approve"
)
# 设置中断策略
workflow.set_interrupt(
"human_approval",
when=lambda x: x.get("require_approval", False)
)
4. 性能优化技巧
4.1 批量中断处理
高频中断场景下,建议实现批量处理机制:
python复制def batch_interrupt(requests: List[InterruptRequest]):
"""批量处理中断请求"""
with get_db_session() as session:
session.bulk_save_mappings(
AuditLog,
[{
"thread_id": r.thread_id,
"action": r.action,
"created_at": datetime.now()
} for r in requests]
)
return [{"status": "pending"} for _ in requests]
4.2 超时自动处理
通过装饰器实现自动超时逻辑:
python复制from functools import wraps
from datetime import timedelta
def timeout(default_action: str, timeout: timedelta = timedelta(hours=24)):
def decorator(f):
@wraps(f)
def wrapper(*args, **kwargs):
start = datetime.now()
result = f(*args, **kwargs)
if datetime.now() - start > timeout:
return Command(resume={"type": default_action})
return result
return wrapper
return decorator
@timeout(default_action="accept")
def approve_large_order(order_id: str):
"""大额订单人工审核"""
...
5. 生产环境注意事项
-
检查点存储:
- 内存存储仅适合开发
- 生产环境需要配置持久化存储
- 定期清理过期检查点
-
中断粒度控制:
- 避免在循环内部使用中断
- 关键路径设置超时机制
- 记录完整审计日志
-
性能监控指标:
python复制PROMETHEUS_COUNTER = Counter( "workflow_interrupts_total", "Total workflow interrupts", ["workflow_name", "node_name"] ) def monitored_interrupt(message: str): PROMETHEUS_COUNTER.labels( current_workflow(), current_node() ).inc() return interrupt(message)
实际部署中发现,当人工审核响应时间超过2小时,工作流恢复成功率会下降30%。建议:
- 实现断点续传机制
- 增加心跳检测
- 对长时间挂起的工作流进行特殊标记
