Nodes、Edges、Shared State、Failure Routing——拆开来看,就是四个能各自单独验证的部分
引言
一张能跑的图,拆开来看,就是四个能各自单独验证的部分。搞清楚这四个构件,大部分困惑就没了。
Nodes(节点)
什么是节点?
一个节点 = 一个跑着自己Loop的Agent或确定性步骤。
节点只认一件事,也只对一件事负责。
节点判断不出”完了”,就不是节点,是隐藏依赖。
节点契约
一个你没法推理的节点,就没法拿去并行。解决办法是给它定一份契约:
`python
# 节点契约示例
def extract_entities(state: GraphState) -> dict:
“””
输入: state.content (str)
输出: state.entities (list[Entity])
失败: 抛出异常,由Failure Routing处理
“””
result = llm.invoke(
f”从以下文本中提取所有实体:n{state[‘content’]}”,
response_format=EntitySchema # 强制结构化输出
)
return {“entities”: result.entities}
`
契约三要素:
- 输入有边界:显式传入,不能指望从共享状态里蹭到
- 输出有形状:最好能用JSON Schema校验
- 失败有声明:做完了就说”完了”,做不了就说”失败”
节点类型
| 类型 | 说明 | 示例 |
|---|
|——|——|——|
| Agent节点 | 调用LLM | 实体提取、内容分析 |
|---|---|---|
| 工具节点 | 调用外部工具 | 文件转换、API调用 |
| 代码节点 | 纯逻辑处理 | 数据过滤、格式转换 |
| 路由节点 | 决定下一步 | 文件类型分类、风险评估 |
Edges(边)
什么是边?
边不只是”B排在A后面”,它是一个关于”传的是什么”的承诺。
按数据给边命名,而不是按顺序命名。
三种边
顺序边:永远触发
`python
graph.add_edge(“scan”, “identify”)
`
条件边:看检查结果
`python
graph.add_conditional_edges(
“classify”,
route_by_type,
{“document”: “extract_doc”, “image”: “extract_image”}
)
`
并行边:一次分给多个
`python
# 从dispatch同时指向verify_a、verify_b、verify_c
graph.add_edge(“dispatch”, “verify_a”)
graph.add_edge(“dispatch”, “verify_b”)
graph.add_edge(“dispatch”, “verify_c”)
`
边的本质
最容易犯的错,是把”然后”当成边。
“总结这个文件,然后告诉我天气”
这两步之间没有边。天气根本用不到那份总结。这其实是两个互不相干的节点,被一段线性脚本硬凑成了先后顺序。
没真用上数据,就没有边。
Shared State(共享状态)
什么是共享状态?
所有节点都要读写的那份数据。这是图的”血液”。
下游节点要用的字段,必须有上游节点写进去。这个对象会逼你承认:这活儿里到底还有多少环节没被真正想清楚。
设计原则
`python
from typing import TypedDict, Annotated
from operator import add
class GraphState(TypedDict):
# 基础字段
file_path: str
file_type: str
content: str
# 节点输出(用add操作符支持追加)
entities: Annotated[list, add] # 多个节点可能都产出实体
relations: Annotated[list, add] # 多个节点可能都产出关系
errors: Annotated[list, add] # 错误日志累加
# 控制字段
status: str
retry_count: int
consecutive_empty: int
`
设计要点:
- 字段名要语义化,不要用data1、data2
- 用Annotated[list, add]支持并行节点追加
- 控制字段和数据字段分开
- 每个字段都要有明确的写入者和读取者
Failure Routing(失败路由)
为什么需要失败路由?
一个节点的重试耗尽了,控制权去哪?
- 退回上一步?
- 转给备用节点?
- 转人工?
没有失败边的图,只是一张流程图,不是一个能跑的系统。
三种失败处理模式
重试
`python
def should_retry(state: GraphState) -> str:
if state[“retry_count”] < 2:
return “retry”
return “fallback”
`
降级
`python
def handle_failure(state: GraphState) -> str:
if state[“has_fallback”]:
return “fallback_node” # 走备用节点
else:
return “mark_failed” # 标记失败,继续其他任务
`
隔离
`python
# parallel()里一个抛错的函数会被解析成null
# 八个正常的agent照样能返回结果
results = await parallel(tasks)
valid = [r for r in results if r is not None] # .filter(Boolean)
`
失败路由设计清单
- [ ] 每个节点都有明确的失败声明
- [ ] 重试次数有上限(建议2次)
- [ ] 重试耗尽后有明确的去向
- [ ] 并行节点失败不影响其他节点
- [ ] 汇入节点能容忍缺失输入
实战:知识大脑的四构件实现
`python
# 1. Nodes
nodes = {
“scan”: scan_directory,
“identify”: identify_file_type, # Magika
“dedup”: check_duplicates,
“extract”: extract_content,
“entities”: extract_entities,
“relations”: extract_relations,
“store”: store_to_database,
}
# 2. Edges
edges = [
(“scan”, “identify”), # 顺序边
(“identify”, “dedup”),
(“dedup”, “extract”),
(“dedup”, “entities”), # 并行边
(“extract”, “relations”),
(“entities”, “relations”),
(“relations”, “store”),
]
# 3. Shared State
class KnowledgeGraphState(TypedDict):
scan_path: str
file_list: list
typed_files: list
new_files: list
content: str
entities: Annotated[list, add]
relations: Annotated[list, add]
stored_ids: list
# 4. Failure Routing
def handle_extract_failure(state):
if state[“retry_count”] < 2:
return “retry_extract”
return “mark_skip” # 跳过这个文件,继续处理其他文件
`
总结
| 构件 | 职责 | 设计要点 |
|---|
|——|——|———-|
| Nodes | 工作单元 | 契约明确,只干一件事 |
|---|---|---|
| Edges | 数据依赖 | 按数据命名,不按顺序 |
| Shared State | 共享数据 | 语义化字段,支持追加 |
| Failure Routing | 失败退路 | 重试+降级+隔离 |
下一章预告:五种基础拓扑模式——链式、扇出、菱形、条件路由、循环。
发表回复