AI 工程技能体系 · 第五篇
– Prompt Engineering — 提示词工程
– Context Engineering — 上下文工程
– Harness Engineering — 工具链工程
– Loop Engineering — 迭代循环工程
– Graph Engineering ← 你在这里
写在前面
Loop Engineering 解决了「怎么让一个 Agent 反复做对一件事」。但当任务复杂到需要多个角色协作、多条路径并行、失败后自动切换路线时,单个 Loop 就不够用了。
Graph Engineering 就是把多个 Loop 组织成一张有向图——节点是工作单元,边是数据依赖,共享状态是大家都要读写的那份数据。
你将收获:
- 理解 Graph 的四个核心构件
- 掌握五种基础拓扑模式
- 用 LangGraph / CrewAI / 原生 Python 实现图编排
- 知道什么时候该用图、什么时候不该用
第一章:先想清楚——你真的需要 Graph 吗?
四个判断条件
| 条件 | 满足 | 不满足 |
|---|
|——|——|——–|
| 任务能拆成不同角色? | ✅ 可以建图 | ❌ 继续用 Loop |
|---|---|---|
| 有真正能并行的子任务? | ✅ 扇出有价值 | ❌ 串行就够了 |
| 单个 Agent 上下文装不下全部背景? | ✅ 拆分释放空间 | ❌ 不用拆 |
| 失败后能负担跳转分支的成本? | ✅ 有失败路由 | ❌ 别建图 |
核心原则:先把一个 Loop 跑稳,再考虑建图。图是 Loop 的组织方式,不是替代品。
Graph vs Loop vs Chain
`
Chain: A → B → C → D (线性,最简单)
Loop: A → B → C → A → B → C (循环,自我改进)
Graph: A → [B,C,D] → E (分支+汇合,并行协作)
↘ F ↗(条件路由)
`
第二章:Graph 的四个核心构件
2.1 Nodes(节点)
一个节点 = 一个跑着自己 Loop 的 Agent 或确定性步骤。
节点契约:
- 输入有边界:显式传入,不能指望从共享状态里蹭到
- 输出有形状:最好能用 JSON Schema 校验
- 只干一件事:做完了就说「完了」,做不了就说「失败」
`python
# 节点示例:提取实体
def extract_entities(state: GraphState) -> dict:
“””输入: state.content 输出: state.entities”””
result = llm.invoke(
f”从以下文本中提取所有实体:n{state[‘content’]}”,
response_format=EntitySchema # 强制结构化输出
)
return {“entities”: result.entities}
`
2.2 Edges(边)
边不只是「B 排在 A 后面」,它是关于「传的是什么」的承诺。
三种边:
| 类型 | 触发条件 | 用途 |
|---|
|——|———-|——|
| 顺序边 | 永远触发 | A 完成后必须执行 B |
|---|---|---|
| 条件边 | 看检查结果 | 分类后走不同路径 |
| 并行边 | 一次分给多个 | 同时执行 B、C、D |
`python
# 条件边示例:按文件类型路由
def route_by_type(state: GraphState) -> str:
file_type = state[“detected_type”]
if file_type == “document”:
return “extract_doc”
elif file_type == “image”:
return “extract_image”
elif file_type == “video”:
return “extract_video”
return “skip”
`
2.3 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
`
2.4 Failure Routing(失败路由)
一个节点失败了,控制权去哪?
`python
# 失败路由示例
def handle_failure(state: GraphState) -> str:
if state[“retry_count”] < 2:
return “retry” # 重试
elif state[“has_fallback”]:
return “fallback_node” # 走备用节点
else:
return “mark_failed” # 标记失败,继续其他任务
`
第三章:五种基础拓扑模式
3.1 链式(Chain)
`
A → B → C → D
`
最简单,每个节点只有一条边进、一条边出。适合线性流程。
`python
from langgraph.graph import StateGraph
graph = StateGraph(GraphState)
graph.add_node(“scan”, scan_file)
graph.add_node(“extract”, extract_content)
graph.add_node(“analyze”, analyze_entities)
graph.add_node(“store”, store_to_db)
graph.add_edge(“scan”, “extract”)
graph.add_edge(“extract”, “analyze”)
graph.add_edge(“analyze”, “store”)
`
3.2 扇出(Fan-out)
`
┌→ B₁ ─┐
A →├→ B₂ ─┤→ C
└→ B₃ ─┘
`
把活儿一次性派出去并行执行。这是图比 Loop 快的核心原因。
`python
# 扇出:并行提取多个文件
from langgraph.graph import END
def fan_out_extract(state: GraphState):
“””一个节点内部并行处理”””
files = state[“pending_files”]
# 并行调用
with ThreadPoolExecutor(max_workers=4) as executor:
futures = {
executor.submit(extract_single, f): f
for f in files
}
results = []
for future in as_completed(futures):
try:
results.append(future.result())
except Exception as e:
results.append({“error”: str(e)})
return {“extracted”: results}
`
3.3 菱形(Diamond)
`
┌→ B ─┐
A →┤ ├→ D
└→ C ─┘
`
派发 → 归约 → 合成。这是最常用的高级拓扑。
`python
# 菱形:并行验证 → 合并结果
graph.add_node(“dispatch”, dispatch_tasks)
graph.add_node(“verify_a”, verify_source_a)
graph.add_node(“verify_b”, verify_source_b)
graph.add_node(“verify_c”, verify_source_c)
graph.add_node(“merge”, merge_results)
graph.add_edge(“dispatch”, “verify_a”)
graph.add_edge(“dispatch”, “verify_b”)
graph.add_edge(“dispatch”, “verify_c”)
graph.add_edge(“verify_a”, “merge”)
graph.add_edge(“verify_b”, “merge”)
graph.add_edge(“verify_c”, “merge”)
`
3.4 条件路由(Conditional Routing)
`
A → [判断] → B(路径1)
→ C(路径2)
`
`python
# 条件路由
graph.add_conditional_edges(
“classify”, # 路由节点
route_by_type, # 路由函数
{
“document”: “extract_doc”, # 路径映射
“image”: “extract_image”,
“video”: “extract_video”,
“skip”: END,
}
)
`
3.5 循环(Loop with Convergence)
`
A → B → [检查] → A(继续)
→ C(完成)
`
关键:必须让它收敛。 用「连续 K 轮无新发现」作为退出条件。
`python
def should_continue(state: GraphState) -> str:
# 连续2轮没有新发现就退出
if state[“consecutive_empty”] >= 2:
return “finalize”
if state[“total_iterations”] >= 10: # 硬上限
return “finalize”
return “continue_loop”
graph.add_conditional_edges(
“discover”,
should_continue,
{“continue_loop”: “discover”, “finalize”: “merge”}
)
`
第四章:框架选型
四大框架对比
| 特性 | LangGraph | CrewAI | AutoGen/MAF | 原生 Python |
|---|
|——|———–|——–|————-|————-|
| 定位 | 底层图编排 | 高层角色编排 | 多Agent对话 | 完全自定义 |
|---|---|---|---|---|
| 学习曲线 | 中等 | 低 | 中等 | 低 |
| 灵活性 | ⭐⭐⭐⭐⭐ | ⭐⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐⭐⭐⭐ |
| 状态管理 | 内置 | 简单 | 内置 | 手动 |
| 持久化 | ✅ 内置 | ❌ | ✅ | 手动 |
| 人机协作 | ✅ 内置 | ❌ | ✅ | 手动 |
| 适合场景 | 复杂生产系统 | 快速原型 | 对话式协作 | 学习/定制 |
LangGraph 示例:知识库文件处理图
`python
from langgraph.graph import StateGraph, END
from typing import TypedDict, Annotated
from operator import add
class FileProcessState(TypedDict):
file_path: str
file_type: str
content: str
entities: Annotated[list, add]
relations: Annotated[list, add]
status: str
# 创建图
graph = StateGraph(FileProcessState)
# 添加节点
graph.add_node(“identify”, identify_file_type) # Magika识别
graph.add_node(“extract_doc”, extract_document) # 文档提取
graph.add_node(“extract_media”, extract_media) # 媒体提取
graph.add_node(“extract_code”, extract_code) # 代码提取
graph.add_node(“analyze”, analyze_entities) # 实体分析
graph.add_node(“store”, store_to_db) # 存储
# 入口
graph.set_entry_point(“identify”)
# 条件路由
graph.add_conditional_edges(
“identify”,
lambda s: s[“file_type”],
{
“document”: “extract_doc”,
“image”: “extract_media”,
“video”: “extract_media”,
“code”: “extract_code”,
}
)
# 汇合到分析
graph.add_edge(“extract_doc”, “analyze”)
graph.add_edge(“extract_media”, “analyze”)
graph.add_edge(“extract_code”, “analyze”)
graph.add_edge(“analyze”, “store”)
graph.add_edge(“store”, END)
# 编译运行
app = graph.compile()
result = app.invoke({“file_path”: “/path/to/file.docx”})
`
CrewAI 示例:研究报告生成
`python
from crewai import Agent, Task, Crew
# 定义角色
researcher = Agent(
role=”研究员”,
goal=”收集全面的研究资料”,
backstory=”你是一位资深研究员,擅长多源信息收集”,
tools=[search_tool, browse_tool]
)
analyst = Agent(
role=”分析师”,
goal=”分析数据并发现洞察”,
backstory=”你是一位数据分析师,擅长发现模式和趋势”
)
writer = Agent(
role=”撰稿人”,
goal=”写出专业的研究报告”,
backstory=”你是一位技术作家,擅长将复杂信息转化为清晰的报告”
)
# 定义任务(CrewAI自动处理依赖和执行顺序)
research_task = Task(
description=”研究 Graph Engineering 的最新进展”,
expected_output=”详细的研究笔记,包含关键发现和来源”,
agent=researcher
)
analysis_task = Task(
description=”分析研究结果,提取关键洞察”,
expected_output=”分析报告,包含趋势和建议”,
agent=analyst,
context=[research_task] # 显式依赖
)
writing_task = Task(
description=”基于分析结果撰写最终报告”,
expected_output=”结构化的研究报告”,
agent=writer,
context=[analysis_task]
)
# 运行
crew = Crew(agents=[researcher, analyst, writer],
tasks=[research_task, analysis_task, writing_task])
result = crew.kickoff()
`
第五章:实战——用原生 Python 实现知识大脑的文件处理图
不用任何框架,纯 Python 实现一个完整的 Graph:
`python
“””
知识大脑 · 文件处理图
拓扑: 扫描 → [扇出] 并行识别/去重 → [并行] 提取/实体/关系 → [汇入] 合并 → 存储
“””
import asyncio
from dataclasses import dataclass, field
from typing import Any, Callable
@dataclass
class Node:
name: str
func: Callable
inputs: list[str] = field(default_factory=list)
outputs: list[str] = field(default_factory=list)
@dataclass
class Edge:
from_node: str
to_node: str
condition: Callable = None # None = 无条件
class GraphEngine:
def __init__(self):
self.nodes: dict[str, Node] = {}
self.edges: list[Edge] = []
self.state: dict[str, Any] = {}
def add_node(self, name: str, func: Callable, inputs=None, outputs=None):
self.nodes[name] = Node(name, func, inputs or [], outputs or [])
def add_edge(self, from_node: str, to_node: str, condition=None):
self.edges.append(Edge(from_node, to_node, condition))
async def run(self, initial_state: dict) -> dict:
self.state = initial_state
completed = set()
while len(completed) < len(self.nodes):
# 找到所有可以执行的节点(前置节点都已完成)
ready = []
for name, node in self.nodes.items():
if name in completed:
continue
deps = [e.from_node for e in self.edges if e.to_node == name]
if all(d in completed for d in deps):
# 检查条件边
can_run = True
for e in self.edges:
if e.to_node == name and e.condition:
if not e.condition(self.state):
can_run = False
break
if can_run:
ready.append(name)
if not ready:
break # 死锁或完成
# 并行执行所有就绪节点
tasks = []
for name in ready:
node = self.nodes[name]
inputs = {k: self.state.get(k) for k in node.inputs}
tasks.append(self._run_node(name, node.func, inputs))
results = await asyncio.gather(*tasks, return_exceptions=True)
for name, result in zip(ready, results):
if isinstance(result, Exception):
self.state[f”error_{name}”] = str(result)
elif isinstance(result, dict):
self.state.update(result)
completed.add(name)
self.state[“completed_nodes”] = list(completed)
return self.state
async def _run_node(self, name, func, inputs):
if asyncio.iscoroutinefunction(func):
return await func(**inputs)
return func(**inputs)
# ===== 构建知识大脑文件处理图 =====
async def build_knowledge_graph():
engine = GraphEngine()
# 节点定义
engine.add_node(“scan”, scan_directory,
outputs=[“file_list”])
engine.add_node(“identify”, identify_file_type,
inputs=[“file_list”], outputs=[“typed_files”])
engine.add_node(“dedup”, check_duplicates,
inputs=[“typed_files”], outputs=[“new_files”])
engine.add_node(“extract”, extract_content,
inputs=[“new_files”], outputs=[“content”])
engine.add_node(“entities”, extract_entities,
inputs=[“content”], outputs=[“entities”])
engine.add_node(“relations”, extract_relations,
inputs=[“content”, “entities”], outputs=[“relations”])
engine.add_node(“store”, store_to_database,
inputs=[“content”, “entities”, “relations”],
outputs=[“stored_ids”])
# 边定义
engine.add_edge(“scan”, “identify”)
engine.add_edge(“identify”, “dedup”)
engine.add_edge(“dedup”, “extract”)
engine.add_edge(“dedup”, “entities”) # 扇出:并行
engine.add_edge(“extract”, “relations”) # 需要content
engine.add_edge(“entities”, “relations”) # 需要entities
engine.add_edge(“relations”, “store”)
# 运行
result = await engine.run({
“scan_path”: “/home/climbing/Documents”,
“db_path”: “data/knowledge.db”
})
print(f”处理完成: {len(result.get(‘stored_ids’, []))} 条知识”)
`
第六章:最佳实践与反模式
✅ 最佳实践
- 先跑稳单个 Loop,再建图——图是 Loop 的组织方式
- 给每个节点定契约——输入有边界,输出有形状
- 边按数据命名,不按顺序命名——能一眼看出数据流
- 默认用 Pipeline,只在必须同步时用屏障——屏障会等最慢的节点
- 给不同节点分配不同档位的模型——重复性工作用便宜模型
- 每个汇入点都设计成能容忍缺失输入——一个节点失败不拖垮全图
❌ 反模式
- 没有真实并行任务却建图——加节点只加成本
- 节点太大或太小——太大装不下上下文,太小协调开销超过收益
- 没有失败路由——图只是一张流程图,不是能跑的系统
- 循环不设收敛条件——变成烧钱死循环
- 所有节点用最强模型——浪费 token
总结
`
Loop Engineering: 一个 Agent 怎么反复做对一件事
Graph Engineering: 多个 Loop 怎么组织成一张能自己路由的图
`
核心公式:
`
Graph = Nodes(工作单元)+ Edges(数据依赖)+ Shared State(共享数据)+ Failure Routing(失败退路)
`
什么时候用 Graph? 当你需要:
- 多个角色协作
- 多条路径并行
- 失败后自动切换
- 超出单个 Agent 上下文限制
什么时候不用? 当你的 Loop 还没跑稳。
*下一篇:Observability Engineering — 可观测性工程*
发表回复