AI PRO·Graph Day 5:Graph Engineering 实操手册

作者:

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’, []))} 条知识”)

`


第六章:最佳实践与反模式

✅ 最佳实践

  1. 先跑稳单个 Loop,再建图——图是 Loop 的组织方式
  2. 给每个节点定契约——输入有边界,输出有形状
  3. 边按数据命名,不按顺序命名——能一眼看出数据流
  4. 默认用 Pipeline,只在必须同步时用屏障——屏障会等最慢的节点
  5. 给不同节点分配不同档位的模型——重复性工作用便宜模型
  6. 每个汇入点都设计成能容忍缺失输入——一个节点失败不拖垮全图

❌ 反模式

  1. 没有真实并行任务却建图——加节点只加成本
  2. 节点太大或太小——太大装不下上下文,太小协调开销超过收益
  3. 没有失败路由——图只是一张流程图,不是能跑的系统
  4. 循环不设收敛条件——变成烧钱死循环
  5. 所有节点用最强模型——浪费 token

总结

`

Loop Engineering: 一个 Agent 怎么反复做对一件事

Graph Engineering: 多个 Loop 怎么组织成一张能自己路由的图

`

核心公式:

`

Graph = Nodes(工作单元)+ Edges(数据依赖)+ Shared State(共享数据)+ Failure Routing(失败退路)

`

什么时候用 Graph? 当你需要:

  • 多个角色协作
  • 多条路径并行
  • 失败后自动切换
  • 超出单个 Agent 上下文限制

什么时候不用? 当你的 Loop 还没跑稳。


*下一篇:Observability Engineering — 可观测性工程*

评论

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注