1. 项目概述:这不是写个脚本,而是给业务流程装上AI引擎
“Python + Agent SDK:把业务场景转化为自动化工作流的完整流程”——这个标题里藏着一个被很多人忽略的关键动词:“转化”。它不是教你怎么调用一个API,也不是演示如何让大模型吐出一段漂亮文案,而是直指企业级落地最痛的环节:如何把模糊的、带上下文的、有判断有回退的业务需求,稳稳地塞进代码可执行、状态可追踪、错误可恢复的自动化流水线里。我做过二十多个从0到1的AI工作流项目,从电商客服工单自动分派,到金融风控报告生成,再到制造业设备报修单的多轮信息补全,踩过的坑比读过的文档还多。核心经验就一条:90%的失败,不是卡在模型能力上,而是卡在“业务逻辑”和“AI行为”之间那层薄薄却极难穿透的膜上。Agent SDK(尤其是LangGraph)的价值,正在于它提供了这层膜的“显微镜”和“缝合针”。它不替代你思考业务规则,但强制你把“如果用户说‘我要退货’,先查订单状态,再判断是否超时,超时则触发人工审核,否则走自动退款”这种自然语言描述,拆解成可图化、可版本化、可压测的节点与边。Python在这里不是配角,而是整个系统的骨架——它负责连接数据库、调用内部API、处理二进制文件、做数值计算,而Agent SDK只管调度和决策。所以,这不是一个“学完就能用”的速成课,而是一套“业务翻译员”的上岗培训:教你如何把老板一句“让销售线索自动跟进”这种模糊指令,翻译成机器能一丝不苟执行的、带心跳检测、带重试策略、带人工兜底开关的生产级工作流。适合谁?不是纯算法工程师,也不是只会写CRUD的后端,而是那些天天泡在业务方会议室、听得懂“商机池”和“SOP”的技术骨干,或者正被重复性报表、跨系统数据搬运压得喘不过气的业务分析师。接下来的内容,没有一行是凭空编造的理论,每一个步骤、每一个参数、每一个报错提示,都来自我上周刚上线的“合同关键条款智能提取与合规初筛”项目现场。
2. 核心设计思路:为什么必须放弃“链式调用”,拥抱“图状工作流”
2.1 传统LangChain Chain的致命短板:它天生不适合真实业务
很多刚接触Agent SDK的人,第一反应是“用LangChain Chain不就行了?”。我试过,也劝退过三个客户。Chain的本质是线性管道:Input → Step1 → Step2 → ... → Output。它像一条笔直的高速公路,所有车(数据)都按固定顺序跑。但真实业务是什么?是立交桥。比如处理一个采购申请单:第一步是OCR识别PDF,第二步是提取金额和供应商,第三步要并行做两件事——查ERP里该供应商的信用额度,同时查财务系统里本月预算余额。如果信用额度不足,要发邮件给采购经理;如果预算余额不足,要触发预算调整流程;如果都OK,才走审批流。这个“并行+条件分支+外部系统交互+人工介入点”的结构,Chain根本无法优雅表达。你硬要用Chain,最后会写出一堆if-else嵌套在RunnableLambda里,调试时看日志像在读天书,加个新分支就得重构整个链。更麻烦的是错误处理——Chain里一个节点挂了,整条链就断了,你想让它自动重试三次再告警?得自己手写状态机。这已经不是“用工具”,而是在“造工具”。
2.2 LangGraph的破局点:用有向无环图(DAG)建模业务逻辑
LangGraph的核心思想,是把工作流当成一张图来画。节点(Node)是具体的执行单元,比如“调用OCR API”、“查询ERP数据库”、“生成合规检查报告”;边(Edge)是节点间的流转规则,比如“OCR成功后,去提取字段”、“提取字段失败,则跳转到人工复核节点”。这张图是有向无环图(DAG),意味着它天然支持:
- 并行执行:你可以定义一个
start_node,它同时触发ocr_task和erp_check_task两个节点,它们互不等待; - 条件路由:
ocr_task输出一个字典,里面包含status: "success"或"failed",边上的condition函数根据这个值决定下一站去哪; - 状态持久化:整个图的运行状态(每个节点的输入/输出、当前在哪个节点、重试次数)都存在一个
State对象里,这个对象可以是简单的dict,也可以是继承自BaseModel的复杂类,方便你存任意业务数据; - 循环与中断:需要人工审核?加个
human_review_node,它输出{"action": "approve"}或{"action": "reject"},边上的条件判断后,要么结束,要么回到rework_node重新处理。
我拿“合同条款提取”项目举例。原始需求是:“上传PDF合同,自动识别甲方、乙方、总金额、付款周期、违约金条款,并对违约金是否超过法定上限做初步判断,有问题标红并通知法务”。用LangGraph,我画出了这张图:
[Upload PDF] ↓ (触发) [Run OCR] → [Extract Entities] → [Validate Penalty Clause] → [Generate Report] ↓ (OCR失败) ↑ (验证失败) ↓ (报告生成失败) [Alert Human] ←───────── [Human Review] ←─────────────── [Retry Logic]注意那个↑箭头,它表示“验证失败”这个状态会反向触发Human Review节点,而不是像Chain那样只能往前走。这个反向、多路径、带状态的特性,才是支撑复杂业务的底层能力。选择LangGraph不是因为它“新”,而是因为它的DAG模型,和我们画业务流程图(BPMN)的思维完全同构。你不需要说服业务方学编程,你直接把他们熟悉的流程图,用Python代码“画”出来,他们一眼就能看懂、能提意见、能确认。
2.3 Python的角色再定位:从胶水语言升级为工作流操作系统
很多人低估了Python在这个架构里的分量。它绝不仅仅是调用几个SDK的“胶水”。在LangGraph里,Python承担着三重操作系统级职责:
- 状态管理器(State Manager):
State类不是随便定义的。我定义的ContractState长这样:
from typing import List, Optional, Dict, Any from pydantic import BaseModel class ContractState(BaseModel): pdf_path: str # 原始文件路径 ocr_text: Optional[str] = None # OCR结果 entities: Dict[str, str] = {} # 提取的实体:{"party_a": "XX公司", "amount": "500000"} penalty_valid: bool = True # 违约金是否合规 penalty_reason: str = "" # 不合规原因 human_review_needed: bool = False # 是否需人工 retry_count: int = 0 # 当前重试次数 audit_log: List[str] = [] # 全流程操作日志,用于审计这个类就是整个工作流的“内存”。每个节点函数的签名都是def node_func(state: ContractState) -> ContractState,它接收当前全部状态,处理后返回一个可能修改了部分字段的新状态。Python的动态类型和Pydantic的强校验,保证了状态流转既灵活又安全。 2.外部系统网关(Gateway):所有和外部世界的交互,都由Python封装。比如调用OCR服务,我不会在节点里直接写requests.post(),而是写一个OcrClient类:
class OcrClient: def __init__(self, api_url: str, timeout: int = 30): self.session = requests.Session() self.session.headers.update({"Authorization": f"Bearer {os.getenv('OCR_TOKEN')}"}) self.api_url = api_url self.timeout = timeout def run_ocr(self, pdf_path: str) -> str: with open(pdf_path, "rb") as f: files = {"file": f} resp = self.session.post(f"{self.api_url}/v1/ocr", files=files, timeout=self.timeout) resp.raise_for_status() return resp.json()["text"]这样,节点函数就变得极其干净:
def ocr_node(state: ContractState) -> ContractState: try: client = OcrClient("https://api.our-ocr.com") text = client.run_ocr(state.pdf_path) state.ocr_text = text state.audit_log.append(f"OCR completed for {state.pdf_path}") except Exception as e: state.audit_log.append(f"OCR failed: {str(e)}") state.human_review_needed = True return state- 可观测性中枢(Observability Hub):工作流跑起来后,你得知道它在哪卡住了。Python利用
langgraph.checkpoint.sqlite.SqLiteSaver把每一步状态存进SQLite,再配合langgraph.prebuilt.tool_node的on_tool_start/on_tool_end钩子,我能精确记录:2024-06-15 14:22:03 | [Extract Entities] | STARTED | input: {'text': '甲方:ABC公司...'} | node_id: extract_001。这些日志不是为了炫技,而是当业务方问“为什么这份合同没生成报告?”时,我打开数据库,5秒内就能定位到是Validate Penalty Clause节点里一个正则表达式没匹配到“违约金”字样,导致penalty_valid始终是True,流程就跳过了人工审核。这种级别的可观测性,是任何黑盒SDK都无法提供的。
3. 实操全流程:从零搭建一个可运行、可调试、可上线的合同工作流
3.1 环境准备与依赖锁定:别让版本冲突毁掉三天工作
别跳过这一步。我见过太多人卡在pip install langgraph之后,发现langchain版本和pydantic冲突,最后在Stack Overflow上浪费一整天。我的方案是:用pip-tools做依赖锁定,用venv隔离环境,用pre-commit保证代码风格。具体命令如下:
# 1. 创建干净虚拟环境 python -m venv ./venv_contract source ./venv_contract/bin/activate # Linux/Mac # ./venv_contract/Scripts/activate # Windows # 2. 安装pip-tools(用于依赖管理) pip install pip-tools # 3. 编写requirements.in,只写顶层依赖 echo "langgraph>=0.1.0" > requirements.in echo "langchain>=0.1.0" >> requirements.in echo "pydantic>=2.0.0,<3.0.0" >> requirements.in echo "requests>=2.28.0" >> requirements.in echo "SQLAlchemy>=2.0.0" >> requirements.in echo "pandas>=2.0.0" >> requirements.in # 4. 生成锁定文件(会解析所有传递依赖,确保可重现) pip-compile requirements.in --output-file requirements.txt # 5. 安装锁定后的依赖 pip install -r requirements.txt这个requirements.txt文件,就是你的“环境DNA”。把它和代码一起提交到Git,任何人git clone后,只要pip install -r requirements.txt,就能得到和你开发时完全一致的环境。特别注意pydantic的版本锁死在<3.0.0,因为LangGraph 0.1.x目前深度依赖Pydantic v2的BaseModel行为,v3的@model_validator语法会破坏状态序列化。这是我在langgraph==0.1.17和pydantic==2.7.1组合下实测稳定的版本对。另外,SQLAlchemy不是LangGraph必需的,但它是SqLiteSaver的底层依赖,必须显式安装,否则checkpoint功能会静默失效。
3.2 定义状态与节点:用代码“画”出你的业务流程图
现在,让我们把前面设想的合同工作流,真正用代码实现。核心是两个文件:state.py定义状态,nodes.py定义节点函数。
state.py:业务状态的宪法
from typing import List, Optional, Dict, Any from pydantic import BaseModel, Field from datetime import datetime class ContractState(BaseModel): # 基础输入 pdf_path: str = Field(..., description="上传的PDF合同文件绝对路径") # OCR阶段 ocr_text: Optional[str] = Field(default=None, description="OCR识别出的纯文本") ocr_error: Optional[str] = Field(default=None, description="OCR错误信息") # 实体提取阶段 entities: Dict[str, str] = Field(default_factory=dict, description="提取的结构化实体") extract_error: Optional[str] = Field(default=None, description="提取错误信息") # 合规验证阶段 penalty_valid: bool = Field(default=True, description="违约金条款是否符合法定要求") penalty_reason: str = Field(default="", description="不合规的具体原因") penalty_clause_text: Optional[str] = Field(default=None, description="识别出的违约金条款原文") # 人工审核与流程控制 human_review_needed: bool = Field(default=False, description="是否需要人工介入") human_review_comment: str = Field(default="", description="人工审核的备注") retry_count: int = Field(default=0, description="当前重试次数,最大3次") # 审计与元数据 start_time: datetime = Field(default_factory=datetime.now, description="工作流启动时间") current_step: str = Field(default="start", description="当前执行的步骤名称") audit_log: List[str] = Field(default_factory=list, description="详细操作日志") # 报告输出 report_html: Optional[str] = Field(default=None, description="最终生成的HTML报告内容") report_path: Optional[str] = Field(default=None, description="报告保存的文件路径") def log(self, message: str): """便捷日志方法""" timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") self.audit_log.append(f"[{timestamp}] {message}") self.log(f"Current step: {self.current_step}")nodes.py:每个节点就是一个独立的、可测试的函数
import re import os import logging from typing import Dict, Any, Optional from state import ContractState from utils.ocr_client import OcrClient # 我们稍后会创建这个 from utils.llm_client import LlmClient # 同样,稍后创建 # 配置日志 logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) def start_node(state: ContractState) -> ContractState: """起始节点:仅更新状态,不做实际工作""" state.current_step = "start" state.log("Workflow started.") return state def ocr_node(state: ContractState) -> ContractState: """OCR节点:调用OCR服务,处理结果""" state.current_step = "ocr" state.log(f"Starting OCR for {state.pdf_path}") try: # 初始化OCR客户端(这里用配置文件或环境变量) ocr_client = OcrClient( api_url=os.getenv("OCR_API_URL", "http://localhost:8000"), timeout=int(os.getenv("OCR_TIMEOUT", "60")) ) text = ocr_client.run_ocr(state.pdf_path) state.ocr_text = text state.log(f"OCR completed. Text length: {len(text)} chars.") except Exception as e: error_msg = f"OCR failed: {str(e)}" state.ocr_error = error_msg state.human_review_needed = True state.log(error_msg) return state def extract_entities_node(state: ContractState) -> ContractState: """实体提取节点:用LLM从OCR文本中提取关键字段""" state.current_step = "extract_entities" state.log("Starting entity extraction...") if not state.ocr_text: state.extract_error = "No OCR text available for extraction." state.human_review_needed = True state.log(state.extract_error) return state try: llm_client = LlmClient( model_name=os.getenv("LLM_MODEL", "gpt-4-turbo"), api_key=os.getenv("OPENAI_API_KEY") ) # 构建提示词(Prompt Engineering是核心!) prompt = f"""你是一个专业的合同审查助手。请从以下合同文本中,精准提取以下字段,以JSON格式输出,不要任何额外解释: - party_a: 甲方全称(必须是公司名,不能是'甲方') - party_b: 乙方全称 - total_amount: 合同总金额(数字,单位:元,去除逗号和'人民币'字样) - payment_terms: 付款周期(如'月结30天'、'验收后付清') - penalty_clause: 违约金条款原文(找到包含'违约金'、'滞纳金'、'赔偿金'等关键词的完整句子) 合同文本: {state.ocr_text[:5000]} # 限制长度,避免超token 输出格式严格为:{{"party_a": "...", "party_b": "...", "total_amount": ..., "payment_terms": "...", "penalty_clause": "..."}}""" response = llm_client.invoke(prompt) # 解析LLM返回的JSON字符串 import json extracted = json.loads(response.strip()) # 数据清洗与校验 if "total_amount" in extracted and isinstance(extracted["total_amount"], str): # 尝试将金额字符串转为数字 try: amount_str = re.sub(r"[^\d.]", "", extracted["total_amount"]) extracted["total_amount"] = float(amount_str) if amount_str else 0.0 except: pass state.entities = extracted state.log(f"Entities extracted: {list(extracted.keys())}") except json.JSONDecodeError as e: state.extract_error = f"LLM output is not valid JSON: {str(e)}" state.human_review_needed = True state.log(state.extract_error) except Exception as e: state.extract_error = f"Entity extraction failed: {str(e)}" state.human_review_needed = True state.log(state.extract_error) return state def validate_penalty_node(state: ContractState) -> ContractState: """违约金验证节点:检查条款是否合法""" state.current_step = "validate_penalty" state.log("Starting penalty clause validation...") if not state.entities or "penalty_clause" not in state.entities: state.penalty_valid = False state.penalty_reason = "No penalty clause found in contract text." state.human_review_needed = True state.log(state.penalty_reason) return state clause_text = state.entities["penalty_clause"].lower() state.penalty_clause_text = clause_text # 法定上限逻辑(简化版,真实项目需对接法律知识库) # 假设:违约金不得超过实际损失的30%,且不得约定为固定高额(如'100万元') # 这里用启发式规则 if "100万" in clause_text or "一百万" in clause_text or "¥1000000" in clause_text: state.penalty_valid = False state.penalty_reason = "Penalty amount is fixed and excessively high (≥1 million RMB)." elif "每日" in clause_text or "每天" in clause_text: # 检查日利率是否过高(例如>0.05%) if re.search(r"(\d+\.?\d*)%.*?日|日.*?(\d+\.?\d*)%", clause_text): # 简单匹配,实际应更严谨 state.penalty_valid = False state.penalty_reason = "Daily penalty rate is potentially illegal." else: # 默认认为合理 state.penalty_valid = True state.penalty_reason = "No obvious violation detected." if not state.penalty_valid: state.human_review_needed = True state.log(f"Penalty validation failed: {state.penalty_reason}") return state def generate_report_node(state: ContractState) -> ContractState: """报告生成节点:生成HTML报告""" state.current_step = "generate_report" state.log("Generating HTML report...") try: from jinja2 import Template template_str = """ <!DOCTYPE html> <html> <head><title>合同合规初筛报告</title> <style>body{font-family:Arial,sans-serif;margin:40px}.header{color:#2c3e50}.entity{margin:10px 0}.warning{color:red;font-weight:bold}</style> </head> <body> <h1 class="header">合同合规初筛报告</h1> <p><strong>文件:</strong>{{ state.pdf_path }}</p> <p><strong>时间:</strong>{{ state.start_time.strftime('%Y-%m-%d %H:%M:%S') }}</p> <h2>提取的实体</h2> {% for key, value in state.entities.items() %} <div class="entity"><strong>{{ key }}:</strong> {{ value }}</div> {% endfor %} <h2>违约金条款验证</h2> <div class="entity"> <strong>条款原文:</strong> {{ state.penalty_clause_text or 'N/A' }} </div> <div class="entity"> <strong>验证结果:</strong> {% if state.penalty_valid %} <span style="color:green">✅ 合规</span> {% else %} <span class="warning">❌ 不合规:{{ state.penalty_reason }}</span> {% endif %} </div> <h2>审计日志</h2> <pre>{{ state.audit_log | join('\n') }}</pre> </body> </html> """ template = Template(template_str) html_content = template.render(state=state) # 保存报告 report_dir = "./reports" os.makedirs(report_dir, exist_ok=True) report_filename = f"report_{int(state.start_time.timestamp())}.html" report_path = os.path.join(report_dir, report_filename) with open(report_path, "w", encoding="utf-8") as f: f.write(html_content) state.report_html = html_content state.report_path = report_path state.log(f"Report generated at {report_path}") except Exception as e: state.log(f"Report generation failed: {str(e)}") # 报告生成失败,不阻断流程,但标记 state.report_path = None return state def human_review_node(state: ContractState) -> ContractState: """人工审核节点:这是一个占位符,实际中会集成到钉钉/企微机器人""" state.current_step = "human_review" state.log("Human review required. Sending notification...") # 这里模拟发送通知(真实项目会调用Webhook) notification_msg = f""" 【合同合规初筛待审】 文件:{state.pdf_path} 问题:{state.penalty_reason or state.ocr_error or state.extract_error} 时间:{state.start_time.strftime('%Y-%m-%d %H:%M:%S')} 查看日志:{state.audit_log[-3:] if state.audit_log else 'No logs'} """ logger.info(notification_msg) # TODO: 调用钉钉机器人Webhook # requests.post("https://oapi.dingtalk.com/robot/send?access_token=xxx", ...) return state看到这里,你应该明白了:每个节点函数,就是一个单一职责、可独立单元测试的Python函数。它只关心自己的输入(state)和输出(修改后的state),不关心其他节点怎么运行。这种设计,让调试变得无比简单——你可以单独运行ocr_node(ContractState(pdf_path="/tmp/test.pdf")),看它是否真的能拿到OCR文本,而不用启动整个工作流。
3.3 构建图与定义边:用Python代码“绘制”你的流程图
有了状态和节点,现在用LangGraph的StateGraph把它们连起来。这是最关键的一步,也是最容易出错的地方。核心是理解add_node和add_edge/add_conditional_edges的区别。
workflow.py:工作流的蓝图
from langgraph.graph import StateGraph, END from langgraph.checkpoint.sqlite import SqliteSaver from nodes import ( start_node, ocr_node, extract_entities_node, validate_penalty_node, generate_report_node, human_review_node ) from state import ContractState # 1. 创建图实例 workflow = StateGraph(ContractState) # 2. 添加所有节点(注意:节点名必须是字符串,且唯一) workflow.add_node("start", start_node) workflow.add_node("ocr", ocr_node) workflow.add_node("extract_entities", extract_entities_node) workflow.add_node("validate_penalty", validate_penalty_node) workflow.add_node("generate_report", generate_report_node) workflow.add_node("human_review", human_review_node) # 3. 定义边(Edge):无条件的、确定性的流转 workflow.add_edge("start", "ocr") workflow.add_edge("ocr", "extract_entities") workflow.add_edge("extract_entities", "validate_penalty") workflow.add_edge("validate_penalty", "generate_report") # 4. 定义条件边(Conditional Edge):根据状态决定下一步 # 这是核心!定义一个函数,它接收state,返回下一个节点名 def route_after_validate(state: ContractState) -> str: """验证后路由:如果需要人工审核,去human_review,否则结束""" if state.human_review_needed: return "human_review" else: return END # END是LangGraph内置的终止节点 # 将条件函数绑定到"validate_penalty"节点的输出上 workflow.add_conditional_edges( "validate_penalty", route_after_validate, { "human_review": "human_review", END: END } ) # 5. 定义另一个条件边:人工审核后的路由 def route_after_human_review(state: ContractState) -> str: """人工审核后路由:假设人工审核后,总是重新走OCR(因为可能需要重传PDF)""" # 真实项目中,这里会根据人工输入的action字段来判断 # 例如:state.human_action == "reprocess" -> "ocr", "approve" -> END return "ocr" # 简化版,循环处理 workflow.add_conditional_edges( "human_review", route_after_human_review, { "ocr": "ocr", END: END } ) # 6. 设置入口点 workflow.set_entry_point("start") # 7. (可选)设置检查点(Checkpoint):让工作流可中断、可恢复 # 使用SQLite作为状态存储后端 memory = SqliteSaver.from_conn_string(":memory:") # 内存数据库,用于演示 # memory = SqliteSaver.from_conn_string("./checkpoints.db") # 生产环境用文件数据库 # 8. 编译图,得到可执行的App app = workflow.compile(checkpointer=memory) # 9. (可选)可视化图结构(需要安装graphviz) # app.get_graph().draw_mermaid_png(output_file_path="contract_workflow.png")这段代码,就是你整个工作流的“源代码级流程图”。add_edge是直线,add_conditional_edges是带标签的分支线。route_after_validate函数就是那个“决策点”,它读取state.human_review_needed这个布尔值,决定是走向human_review还是END。这个函数的返回值,必须是add_conditional_edges第二个参数字典里的key。如果你写错了,比如返回了"review"但字典里是"human_review",工作流就会卡死,没有任何错误提示——这是新手最常见的坑。我建议在route_*函数里加一行logger.debug(f"Routing to: {next_node}"),方便调试。
3.4 运行、调试与监控:让工作流从“能跑”到“稳跑”
现在,工作流编译好了,怎么运行?怎么知道它在哪一步卡住了?怎么模拟人工审核?这才是工程化的关键。
run_workflow.py:一个健壮的运行器
import asyncio import sys from pathlib import Path from workflow import app from state import ContractState async def main(): # 1. 构建初始状态 if len(sys.argv) < 2: print("Usage: python run_workflow.py <path_to_pdf>") sys.exit(1) pdf_path = sys.argv[1] if not Path(pdf_path).exists(): print(f"PDF file not found: {pdf_path}") sys.exit(1) initial_state = ContractState(pdf_path=pdf_path) # 2. 运行工作流(异步) # config={"configurable": {"thread_id": "123"}} 是必须的,用于checkpoint config = {"configurable": {"thread_id": "test_thread_001"}} print(f"Starting workflow for {pdf_path}...") try: # stream() 方法会返回一个生成器,逐个yield出每一步的状态 async for event in app.astream(initial_state, config=config, stream_mode="values"): # event 就是每一步执行后的ContractState print(f"--- Step: {event.current_step} ---") if event.ocr_error: print(f" OCR Error: {event.ocr_error}") if event.extract_error: print(f" Extract Error: {event.extract_error}") if event.penalty_reason: print(f" Penalty Check: {event.penalty_reason}") if event.human_review_needed: print(f" ⚠️ HUMAN REVIEW NEEDED!") # 只打印关键字段,避免刷屏 print(f" Entities: {list(event.entities.keys()) if event.entities else 'None'}") print(f" Report Path: {event.report_path or 'Not generated'}") print() except Exception as e: print(f"Workflow execution failed: {e}") # 打印完整的审计日志,用于排查 print("Full audit log:") for log in initial_state.audit_log: print(f" {log}") sys.exit(1) print("Workflow completed.") if __name__ == "__main__": asyncio.run(main())运行它:python run_workflow.py ./samples/contract_v1.pdf。你会看到实时的、按步骤输出的日志,清晰地告诉你每一步发生了什么。这就是astream的价值——它不是等整个流程跑完才给你结果,而是“流式”输出,让你能实时监控。
调试技巧:
- 断点调试:在VS Code里,直接在
ocr_node函数第一行打个断点,然后运行run_workflow.py,它会在OCR调用前停下来,你可以检查state.pdf_path是否正确,os.getenv("OCR_API_URL")是否已设置。 - 状态快照:
app.get_state(config)可以获取当前线程的最新状态。在run_workflow.py里加一行print(app.get_state(config)),就能看到此刻所有字段的值。 - 重放(Replay):如果工作流在
validate_penalty节点失败了,你不需要重跑整个OCR和提取,只需:
这种能力,在处理耗时的OCR或LLM调用时,能节省大量时间。# 获取失败时的状态 state_at_failure = app.get_state(config) # 修改状态,模拟人工修正 state_at_failure.entities["penalty_clause"] = "违约金为合同总额的10%,不超过实际损失的30%。" # 从validate_penalty节点重新开始 config["configurable"]["thread_ts"] = state_at_failure.values["__end_time__"] # 需要设置时间戳 app.invoke(state_at_failure, config=config)
监控与告警:生产环境必须加监控。我在workflow.py里加了一个全局钩子:
from langgraph.prebuilt import ToolNode from langgraph.graph import START # 在workflow.compile()之前,添加一个钩子 def on_node_start(node_name: str, state: ContractState): logger.info(f"Node '{node_name}' STARTED at {datetime.now()}") def on_node_end(node_name: str, state: ContractState, result: Any): duration = (datetime.now() - state.start_time).total_seconds() logger.info(f"Node '{node_name}' ENDED. Duration: {duration:.2f}s") # 注册钩子(LangGraph 0.1.x 的方式) # 注意:这需要在app.compile()之后,但在使用前 # 实际项目中,我会用Prometheus + Grafana,记录每个节点的P95延迟和错误率这些日志,配合ELK(Elasticsearch, Logstash, Kibana)堆栈,就能构建出一张实时的工作流健康仪表盘:哪个节点最慢?哪个节点错误率最高?今天有多少合同进入了人工审核队列?这才是真正的“让AI下地干活”。
4. 常见问题与独家避坑指南:那些文档里不会写的血泪教训
4.1 “工作流卡在某个节点不动了!”——状态未更新是最大陷阱
现象:你运行app.invoke(),控制台只输出Starting workflow...,然后就没了,程序既不报错也不结束。
根本原因:节点函数没有返回修改后的state,或者返回了错误的state类型。LangGraph的StateGraph要求每个节点函数的签名必须是def node(state: YourState) -> YourState,并且必须返回一个YourState的实例。常见错误有:
- 忘记
return state,函数默认返回None,工作流就卡死了。 - 在节点里做了
state = ContractState(...),创建了一个新对象,但没返回它,而是返回了旧的state(或None)。 - 返回了
dict或其他类型,而不是ContractState的实例。
排查方法:
- 在每个节点函数的末尾,加一行
print(f"Node {node_name} returning state type: {type(state)}")。 - 确保
print语句输出的是<class '__main__.ContractState'>,而不是`<class '