pydantic-graph 全面指南:用类型注解驱动的 Python 图与状态机库
【免费下载链接】pydantic-aiHow Python does AI. Agents, realtime voice, image generation, embeddings. Every model, every interface, typed end to end.项目地址: https://gitcode.com/GitHub_Trending/py/pydantic-ai
导读
本文是 pydantic-graph 的完整技术指南。pydantic-graph 是随 Pydantic AI 一同开发、却完全独立于pydantic-ai的纯图(Graph)与有限状态机(Finite State Machine)库,它允许你用标准 Python 语法——尤其是用节点run方法的返回类型注解来定义节点之间的边(edge),从而以类型安全的方式编排复杂的异步工作流。读完本文,你将掌握:声明式BaseNode与构建器GraphBuilder两种编程范式、图的构建与校验、同步/异步/流式执行方式、条件分支(Decision)、并行分叉与汇聚(Fork/Join/Reducer)、Mermaid 图渲染,以及它如何作为 Pydantic AI Agent 循环的底层引擎工作。
一、pydantic-graph 是什么:一个"无 GenAI 依赖"的图状态机库
pydantic-graph的定位在 pydantic_graph/README.md 中写得很清楚:它作为 Pydantic AI 的一部分被开发,但不依赖pydantic-ai或任何相关包,可以被当作一个纯粹的"基于图的状态机库"来使用——无论你是否在使用 Pydantic AI,甚至是否在做 GenAI 相关开发,它都可能对你有用。
与 Pydantic AI 一脉相承,这个库把类型安全和通用 Python 语法放在首位,刻意回避"晦涩的、领域特有的 Python 语法"(esoteric, domain-specific use of Python syntax)。其核心设计思想是:
- 节点(Node):图中的执行单元,通过
BaseNode子类或普通异步函数定义; - 边(Edge):通过节点的返回类型注解(return type hint)自动推导——这是整个库最与众不同的地方,写图就像写类型标注一样自然;
- 执行引擎:内置基于 anyio 的异步任务调度,支持并行、条件分支、汇聚与取消。
从源码模块的文档字符串可以进一步确认它的地位:pydantic_graph/pydantic_graph/init.py 中写道,它是 "Type-hint based graph librarypowering the Pydantic AI agent loop"——即 Pydantic AI 的 Agent 循环本身就是构建在 pydantic-graph 之上的。因此,理解这个库等于理解了 Pydantic AI 内部编排的骨架。
安装与运行环境
pydantic-graph以独立包发布(见 pydantic_graph/pyproject.toml):
pip install pydantic-graph # 或使用 uv uv add pydantic-graph其运行环境要求与依赖如下:
| 项目 | 要求 |
|---|---|
| Python 版本 | >=3.10(3.10 至 3.14 均受支持) |
anyio | >=4.7.0(取消投递/流式清理依赖,见 pyproject 注释) |
logfire-api | >=3.14.1(OpenTelemetry 可观测性) |
pydantic | >=2.12 |
typing-inspection | >=0.4.0(返回类型注解解析) |
| 许可证 | MIT |
二、第一个图:从 README 的"逢 5 递进"示例说起
README 给出的基础示例非常精炼,它同时演示了声明式节点(BaseNode子类)与构建器(GraphBuilder)两种风格的混用。我们先完整复现它:
from __future__ import annotations from dataclasses import dataclass from pydantic_graph import BaseNode, End, GraphBuilder, GraphRunContext, StepContext @dataclass class DivisibleBy5(BaseNode[None, None, int]): foo: int async def run( self, ctx: GraphRunContext, ) -> Increment | End[int]: if self.foo % 5 == 0: return End(self.foo) else: return Increment(self.foo) @dataclass class Increment(BaseNode): foo: int async def run(self, ctx: GraphRunContext) -> DivisibleBy5: return DivisibleBy5(self.foo + 1) g = GraphBuilder(input_type=int, output_type=int) @g.step async def start(ctx: StepContext[None, None, int]) -> DivisibleBy5: return DivisibleBy5(ctx.inputs) g.add( g.node(DivisibleBy5), g.node(Increment), g.edge_from(g.start_node).to(start), ) fives_graph = g.build() async def main(): result = await fives_graph.run(inputs=4) print(result) #> 5运行main()会得到输出5:从输入4出发,start把输入包装成DivisibleBy5(4),随后在DivisibleBy5与Increment之间循环递增,直到遇到 5 的倍数5时返回End(5)。
这个例子展示了三条核心机制:
- 类型注解即拓扑:
DivisibleBy5.run的返回注解Increment | End[int]告诉图引擎"下一步可以走到Increment节点,或是以int值结束";Increment.run的返回注解DivisibleBy5则指出唯一的后继节点。这些注解在运行时被读取并用于构建边(见下文"运行时的类型反射")。 End是图的出口:任何节点返回End(data)即代表图执行完成,data会成为整个run()的返回值。- 构建器与声明式节点可以混用:
@g.step装饰的普通异步函数start被当作"步骤节点",g.node(DivisibleBy5)把BaseNode子类接入图中,g.edge_from(g.start_node).to(start)显式声明入口边。
运行时的类型反射:边是如何"无中生有"的
"返回类型注解定义边"并非魔法,其底层实现在 pydantic_graph/pydantic_graph/graph_builder.py 的GraphBuilder.node()中:它通过get_type_hints(node_type.run, ...)读取run方法的return注解,再交由_edge_from_return_hint()(graph_builder.py)解析联合类型中的每个成员:
- 成员是
End[...]→ 连接到end_node(结束节点); - 成员是某个
BaseNode子类 → 连接到对应的节点步骤; - 成员是
StepNode[...]/JoinNode[...](带Annotated标注)→ 连接到对应的 step/join。
若run方法缺少返回类型注解,node()会直接抛出GraphSetupError("Node ... is missing a return type hint on itsrunmethod")。同时,basenode.py 的文档字符串明确警告:run方法的返回类型会在运行时被 pydantic-graph 读取,用于定义"接下来可以调用哪些节点",并在图执行时强制执行。
三、核心抽象:BaseNode、End、GraphRunContext 与 Edge
basenode.py 定义了声明式编程范式的四个基石:
BaseNode:节点的抽象基类
class BaseNode(ABC, Generic[StateT, DepsT, NodeRunEndT]): @abstractmethod async def run(self, ctx: GraphRunContext[StateT, DepsT]) -> BaseNode[StateT, DepsT, Any] | End[NodeRunEndT]: ...- 三个类型参数:
StateT(图状态类型,默认object)、DepsT(依赖类型,默认object)、NodeRunEndT(该节点可能的结束值类型,默认Never)。示例中的BaseNode[None, None, int]表示"无状态、无依赖、结束值类型为 int"。 run是唯一必须实现的方法:接收GraphRunContext,返回"下一个节点"或End。返回值中的节点类型就是下一跳,这正是类型驱动拓扑的关键。get_node_id():类方法,默认返回类名(cls.__name__),用作节点在图中唯一 ID;子类可重写。
GraphRunContext:节点执行时的上下文
@dataclass(kw_only=True) class GraphRunContext(Generic[StateT, DepsT]): state: StateT # 图的当前状态 deps: DepsT # 图运行依赖(如数据库句柄、配置对象)state在整个图运行期间共享,是节点间传递可变数据的官方通道;deps用于注入不随运行变化的外部依赖。
End:图的终止信号
@dataclass class End(Generic[RunEndT]): data: RunEndT # 图运行结束后对外返回的数据Edge:边标签
@dataclass(frozen=True) class Edge: label: str | None # 边的标签,主要用于 Mermaid 渲染与可读性Edge是一个标注类型,用来给图中的边打标签,方便调试与可视化。
四、GraphBuilder:构建器的完整 API
除了声明式BaseNode,pydantic-graph 从 1.0 版本起主推构建器 API。GraphBuilder(graph_builder.py)提供了流式(fluent)的建图方式,全部参数与核心方法如下。
构造参数
g = GraphBuilder( name: str | None = None, # 图名称,缺省时在首次调用图方法时从调用帧推断 state_type=..., # 图状态类型,默认 NoneType deps_type=..., # 依赖类型,默认 NoneType input_type=..., # 输入类型,默认 NoneType output_type=..., # 输出类型,默认 NoneType auto_instrument: bool = True, # 是否自动创建 OpenTelemetry instrumentation span )其中state_type/deps_type/input_type/output_type都是TypeOrTypeExpression,即既可以是具体类型,也可以是类型表达式(由util.py的TypeExpression支持)。auto_instrument=True时,每次run都会自动生成形如run graph <name>的 span,每个步骤节点生成run node <id>的 span(见 graph_builder.py 与 graph_builder.py)。
节点与边构建方法一览
| 方法 | 作用 | 关键参数 |
|---|---|---|
start_node/end_node | 获取图的入口/出口节点(属性) | — |
step(call=None, *, node_id=None, label=None) | 把异步函数包装成步骤节点,可作装饰器或直接调用 | node_id缺省取函数名 |
stream(call=None, *, node_id=None, label=None) | 把"返回异步迭代器的函数"包装为流式步骤节点 | 返回值类型为AsyncIterable[OutputT] |
node(node_type) | 将BaseNode子类接入图,自动分析run返回注解建边 | 缺失返回注解抛GraphSetupError |
join(reducer, *, initial/initial_factory, node_id=None, parent_fork_id=None, preferred_parent_fork='farthest') | 创建汇聚节点 | 见"并行与汇聚"一节 |
add(*edges) | 一次性加入多条边路径,自动创建 fork/decision 等中间节点 | 可传EdgePath |
add_edge(source, destination, *, label=None) | 添加一条简单边 | — |
add_mapping_edge(source, map_to, *, pre_map_label, post_map_label, fork_id, downstream_join_id) | 添加"把可迭代数据摊平到并行路径"的边 | downstream_join_id用于空可迭代的兜底 |
edge_from(*sources) | 返回EdgePathBuilder,链式构造路径 | 支持.to()、.label()、.map()、.broadcast()、.transform() |
decision(*, note=None, node_id=None) | 创建条件分支节点 | note会渲染进 Mermaid 图 |
match(source, *, matches=None) | 创建基于类型/谓词的分支匹配器 | 见"条件分支" |
match_node(source, *, matches=None) | 针对BaseNode子类的分支匹配 | — |
build(validate_graph_structure=True) | 完成建图并返回可执行的Graph | 见"构建校验" |
README 示例中g.edge_from(g.start_node).to(start)便是edge_from+to的典型用法:从__start__节点出发,连到start步骤。
build() 内部的四步流水线
build()在返回Graph之前会依次执行(graph_builder.py):
_replace_placeholder_node_ids:为决策/广播等自动生成的节点分配确定性 ID(去重编号);_flatten_paths:把路径中的 Map/Broadcast 标记拆分为显式的 Fork 节点;_normalize_forks:把"有多条出边的普通节点"统一转换为显式广播 fork,保证执行模型单一;_validate_graph_structure(可通过validate_graph_structure=False关闭):校验图中不存在死节点、结束节点必须可达、所有节点必须从 start 可达等。
五、执行一个图:run / run_sync / iter 三种方式
构建完成后得到的是Graph实例(graph_builder.py),它提供三种执行入口:
1.await graph.run(...)—— 异步运行到结束
async def run( self, *, state: StateT = None, deps: DepsT = None, inputs: InputT = None, span: AbstractContextManager[AbstractSpan] | None = None, infer_name: bool = True, ) -> OutputT:- 内部通过
graph.iter()创建GraphRun,循环await graph_run.next(event)直到收到EndMarker,再解包返回最终输出值(graph_builder.py)。 - 源码注释特别提到:
run刻意走next()路径,以保证"逐步入口"在关键路径上被测试到。
2.graph.run_sync(...)—— 同步运行
run_sync是run的同步包装:它在当前事件循环上调用loop.run_until_complete(...),因此不能在已有事件循环的 async 上下文中调用。对不支持run_until_complete的事件循环(如 Temporal 的 workflow 事件循环),会抛出UnsupportedEventLoopError(见 exceptions.py)。
3.graph.iter()+GraphRun—— 逐步执行与错误恢复
async with graph.iter(state=..., deps=..., inputs=...) as graph_run: # graph_run 是 GraphRun 实例 event = await graph_run.next() # 推进一步 graph_run.override_next(new_tasks) # 覆盖下一步(错误恢复/提前结束) event = await graph_run.next(event)GraphRun(graph_builder.py)是单次图执行的实例,暴露以下关键成员:
| 成员 | 说明 |
|---|---|
next(value=None) | 推进一步执行,返回EndMarker或Sequence[GraphTask] |
override_next(value) | 覆盖待执行的下一步:可传新任务序列继续执行,或传EndMarker提前结束 |
next_task | 查看下一步待执行的任务 |
output | 若已结束,返回最终输出 |
state/deps/inputs | 本次运行的上下文 |
| 异步迭代协议 | GraphRun本身是异步迭代器(__aiter__/__anext__) |
错误恢复机制:当某个节点抛出异常时,迭代器不会立即中断,而是产出ErrorMarker(封装原始异常),调用方可在下一次迭代前通过override_next()注入新任务来完成恢复;如果调用方不处理,异常会在下一次__anext__时被重新抛出(graph_builder.py)。ErrorMarker的源码注释说明这正是为after_node_run/on_node_run_error这类钩子系统设计的。
六、步骤节点 Step:函数式编程范式
step.py提供了函数式步骤抽象,让普通异步函数直接参与图执行。
StepContext:步骤的上下文
@dataclass(init=False) class StepContext(Generic[StateT, DepsT, InputT]): # 三个只读属性 @property def state(self) -> StateT: ... # 当前图状态 @property def deps(self) -> DepsT: ... # 图运行依赖 @property def inputs(self) -> InputT: ... # 传入本步骤的输入数据README 示例中的start步骤正是通过ctx.inputs拿到图的初始输入4的。
Step / StepFunction / StreamFunction
StepFunction是异步函数协议:async def f(ctx: StepContext) -> OutputT;StreamFunction是异步迭代器协议:async def f(ctx: StepContext) -> AsyncIterator[OutputT],配合GraphBuilder.stream()使用,可在节点间流式产出数据;Step封装"函数 + 节点 ID + 标签",其as_node(inputs=None)返回StepNode——一个可被BaseNode返回的"指向某步骤"的桥接节点。
StepNode 与 NodeStep:两套范式的桥
StepNode:把 step 和它要接收的输入绑定在一起。BaseNode.run返回StepNode即表示"切换到某个 builder 步骤",它本身不可直接运行(run恒抛NotImplementedError,见 step.py)。NodeStep:反向桥接——把BaseNode子类包装成 builder 步骤,校验输入确为该节点类型后调用其run方法(step.py)。GraphBuilder.node()内部正是创建NodeStep。
这两类桥节点让"声明式节点"与"函数式步骤"可以在同一张图里自由组合,这也是 README 示例能混用的原因。
七、条件分支:Decision 与 match
decision.py实现了基于运行时条件的路由。典型用法:
decision = g.decision(note='根据输入类型路由') decision = decision.branch(g.match(int).to(handle_int)) decision = decision.branch(g.match(str).to(handle_str))分支匹配的默认逻辑
DecisionBranch的matches参数为None时,采用以下默认判定(decision.py 与运行时实现 graph_builder.py):
source类型 | 匹配方式 |
|---|---|
Any/object | 恒为真,该分支总是匹配 |
Literal[...] | 输入值属于字面量集合之一即匹配 |
| 其他类型 | 用isinstance(inputs, source)判定 |
分支按顺序测试,命中第一个匹配分支的路径执行。若所有分支都不匹配,运行时抛出RuntimeError: No branch matched inputs ...。
分支的路径修饰
DecisionBranchBuilder还提供链式方法:
.to(dest, *extra_dests, fork_id=None):指定一个或多个目标(多个目标自动生成广播 fork);.broadcast(get_forks):把分支广播到由回调产生的多个分支路径;.transform(func):在分支路径上施加同步数据转换(TransformFunction接收StepContext返回新值);.map(fork_id=None, downstream_join_id=None):把可迭代输出摊平为多条并行路径;.label(label):为路径段打标签(仅用于 Mermaid 渲染)。
Decision还通过HandledT逆变类型参数在静态类型检查层面保证"所有可能的输入类型都被穷尽处理",_force_handled_contravariant方法(decision.py)是这一穷尽性检查的类型层面实现细节。
八、并行与汇聚:Fork、Map、Join 与 Reducer
并行是 pydantic-graph 的杀手级能力,其机制在 node.py 与 join.py 中定义。
Fork:广播 vs 映射
Fork节点有两种模式(node.py):
is_map | 行为 | 输入要求 |
|---|---|---|
False(广播 broadcast) | 同一份数据被派发到每一条下游路径 | InputT == OutputT |
True(映射 map) | 可迭代输入被摊平,每个元素走一条并行路径 | InputT为Sequence[OutputT] |
downstream_join_id字段用于"映射空可迭代"时的兜底:当映射的输入为空序列时,join 仍需以初始值触发一次(见 graph_builder.py 的注释与实现)。
在构建 API 中,广播通常由.to(dest1, dest2)多目标或PathBuilder.broadcast()隐式创建;映射则由.map()/add_mapping_edge()创建。
Join 与 Reducer:汇聚并行结果
join = g.join( reduce_list_append, # 内置 reducer initial=[], # 或 initial_factory=... parent_fork_id=None, preferred_parent_fork='farthest', # 或 'closest' )Join.reduce()会根据 reducer 的形参个数自动区分两种签名(join.py):
- 普通 reducer:
(current, inputs) -> current - 上下文 reducer:
(ctx: ReducerContext, current, inputs) -> current
内置 reducer 一览(join.py):
| Reducer | 行为 |
|---|---|
reduce_null | 丢弃所有输入,返回None(仅作汇聚信号) |
reduce_list_append | 追加单个元素到列表 |
reduce_list_extend | 用可迭代扩展列表 |
reduce_dict_update | 用映射更新字典 |
reduce_sum | 数值累加 |
ReduceFirstValue | 取第一个到达的值,并取消其余兄弟任务(早期停止) |
ReducerContext暴露state、deps与cancel_sibling_tasks()——后者允许实现"首个结果即返回"的短路语义。ReduceFirstValue正是通过调用它实现取消的。
支配 fork(Dominating Fork):Join 的同步前提
每一个 Join 都必须满足"存在一个支配它的 Fork F":从 StartNode 到该 Join 的所有路径都必须经过 F,且不含绕过 F 的环。构建器通过_collect_dominating_forks(graph_builder.py)自动计算,若找不到支配 fork 会抛出GraphBuildingError并附带 Mermaid 图帮助排查。这一约束保证引擎能判断"该 fork 的所有上游任务是否已完成"。
preferred_parent_fork参数('farthest'/'closest')则在存在多个候选 fork 时决定优先选择哪一个。
九、Mermaid 可视化:一张图看懂你的图
Graph.render()把整个图渲染为 Mermaid 的stateDiagram-v2字符串(graph_builder.py),str(graph)也直接返回渲染结果,因此调试时只需print(graph)。
mermaid_text = fives_graph.render( title='逢 5 递进图', direction='LR', # TB(默认) | LR | RL | BT )渲染细节(graph_builder.py):
- start/end 节点用 Mermaid 的
[*]特殊语法表示; - step 节点输出为
node_id: label; - join 节点输出为
state <id> <<join>>,fork 输出为<<fork>>,decision 输出为<<choice>>并支持note right of <id>附注(对应decision(note=...)); - 边标签(
label/Edge/LabelMarker)会渲染为--> ...: label; - 输出前会做基于 BFS 深度的拓扑排序(
_topological_sort),保证图从 start 向 end 有序展示。
仓库中 tests/graph/builder/test_mermaid_rendering.py 提供了渲染行为的完整测试覆盖。
十、构建校验与异常体系
GraphBuilder.build(validate_graph_structure=True)默认执行结构校验(graph_builder.py),检查项包括:
- start 节点必须存在出边;
- 图中必须存在通向 end 节点的边;
- 除 end 外不允许"死端"节点(无出边、非决策分支);
- end 节点必须从 start 可达;
- 所有节点必须从 start 可达。
若你的图刻意违反上述假设(例如某些节点允许成为死端),可以显式传validate_graph_structure=False跳过校验。校验失败时抛出的GraphValidationError会附带提示语:"If this is intentional, you can suppress this error by passingvalidate_graph_structure=False..."
异常体系定义在 exceptions.py,共五类:
| 异常 | 基类 | 触发场景 |
|---|---|---|
GraphSetupError | TypeError | 图配置错误,如节点run缺失返回注解 |
GraphBuildingError | ValueError | 建图阶段错误,如节点 ID 冲突、Join 缺少支配 fork |
GraphValidationError | ValueError | 图结构校验失败 |
GraphRuntimeError | RuntimeError | 图执行阶段错误 |
UnsupportedEventLoopError | RuntimeError | 同步方法遇到不支持run_until_complete的事件循环 |
十一、pydantic-graph 与 Pydantic AI 的关系
如前所述,pydantic-graph 是 Pydantic AI Agent 循环的底层引擎(见 pydantic_graph/pydantic_graph/init.py 的模块 docstring)。Pydantic AI 的 Agent 内部维护了一张由节点(用户提示节点、模型调用节点、工具执行节点、输出节点等)组成的图,通过pydantic_graph的类型驱动机制完成"调用模型 → 执行工具 → 再调用模型"的循环编排。这意味着:
- 你在 Pydantic AI 中体验到的
Agent.run()的多轮工具调用循环、流式输出、取消与重试,本质都是 pydantic-graph 的图执行能力; - pydantic-graph 完全可以脱离 Pydantic AI 独立使用,用于任何需要"有状态、可并行、可恢复"的异步工作流场景。
仓库的测试目录 tests/graph/builder/ 提供了覆盖各能力的测试参考,包括test_graph_builder.py(建图)、test_graph_execution.py(执行)、test_graph_iteration.py(逐步迭代与恢复)、test_decisions.py(条件分支)、test_joins_and_reducers.py(汇聚)、test_broadcast_and_spread.py(广播与映射)、test_parent_forks.py(支配 fork)、test_mermaid_rendering.py(可视化)等,可作为理解行为边界的权威示例集。
十二、快速上手清单
- 安装:
pip install pydantic-graph(Python ≥ 3.10)。 - 定义节点:写
BaseNode子类,run(ctx: GraphRunContext) -> NextNode | End[OutputT],返回注解即拓扑。 - 建图:
g = GraphBuilder(input_type=..., output_type=...),用@g.step定义入口步骤,用g.add(g.node(...), g.edge_from(g.start_node).to(start))组装。 - 执行:
await g.build().run(inputs=...);需要逐步控制用graph.iter()+GraphRun.next();需要同步用run_sync()。 - 进阶:
g.decision()+g.match()做条件路由;edge_from(...).to(a, b)或.map()做并行;g.join(reducer, initial=...)汇聚并行结果;print(graph)直接获得 Mermaid 图。 - 调试:利用
GraphSetupError/GraphValidationError的报错信息与附带的 Mermaid 图定位结构问题。
【免费下载链接】pydantic-aiHow Python does AI. Agents, realtime voice, image generation, embeddings. Every model, every interface, typed end to end.项目地址: https://gitcode.com/GitHub_Trending/py/pydantic-ai
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考