pydantic-graph 全面指南:用类型注解驱动的 Python 图与状态机库
2026/9/14 12:56:00 网站建设 项目流程

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),随后在DivisibleBy5Increment之间循环递增,直到遇到 5 的倍数5时返回End(5)

这个例子展示了三条核心机制:

  1. 类型注解即拓扑DivisibleBy5.run的返回注解Increment | End[int]告诉图引擎"下一步可以走到Increment节点,或是以int值结束";Increment.run的返回注解DivisibleBy5则指出唯一的后继节点。这些注解在运行时被读取并用于构建边(见下文"运行时的类型反射")。
  2. End是图的出口:任何节点返回End(data)即代表图执行完成,data会成为整个run()的返回值。
  3. 构建器与声明式节点可以混用@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 版本起主推构建器 APIGraphBuilder(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.pyTypeExpression支持)。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):

  1. _replace_placeholder_node_ids:为决策/广播等自动生成的节点分配确定性 ID(去重编号);
  2. _flatten_paths:把路径中的 Map/Broadcast 标记拆分为显式的 Fork 节点;
  3. _normalize_forks:把"有多条出边的普通节点"统一转换为显式广播 fork,保证执行模型单一;
  4. _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_syncrun的同步包装:它在当前事件循环上调用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)推进一步执行,返回EndMarkerSequence[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))

分支匹配的默认逻辑

DecisionBranchmatches参数为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)可迭代输入被摊平,每个元素走一条并行路径InputTSequence[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暴露statedepscancel_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),检查项包括:

  1. start 节点必须存在出边;
  2. 图中必须存在通向 end 节点的边;
  3. 除 end 外不允许"死端"节点(无出边、非决策分支);
  4. end 节点必须从 start 可达;
  5. 所有节点必须从 start 可达。

若你的图刻意违反上述假设(例如某些节点允许成为死端),可以显式传validate_graph_structure=False跳过校验。校验失败时抛出的GraphValidationError会附带提示语:"If this is intentional, you can suppress this error by passingvalidate_graph_structure=False..."

异常体系定义在 exceptions.py,共五类:

异常基类触发场景
GraphSetupErrorTypeError图配置错误,如节点run缺失返回注解
GraphBuildingErrorValueError建图阶段错误,如节点 ID 冲突、Join 缺少支配 fork
GraphValidationErrorValueError图结构校验失败
GraphRuntimeErrorRuntimeError图执行阶段错误
UnsupportedEventLoopErrorRuntimeError同步方法遇到不支持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(可视化)等,可作为理解行为边界的权威示例集。

十二、快速上手清单

  1. 安装pip install pydantic-graph(Python ≥ 3.10)。
  2. 定义节点:写BaseNode子类,run(ctx: GraphRunContext) -> NextNode | End[OutputT],返回注解即拓扑。
  3. 建图g = GraphBuilder(input_type=..., output_type=...),用@g.step定义入口步骤,用g.add(g.node(...), g.edge_from(g.start_node).to(start))组装。
  4. 执行await g.build().run(inputs=...);需要逐步控制用graph.iter()+GraphRun.next();需要同步用run_sync()
  5. 进阶g.decision()+g.match()做条件路由;edge_from(...).to(a, b).map()做并行;g.join(reducer, initial=...)汇聚并行结果;print(graph)直接获得 Mermaid 图。
  6. 调试:利用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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询