ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

pydantic-graph 全面指南:用类型注解驱动的 Python 图与状态机库

pydantic-graph 全面指南:用类型注解驱动的 Python 图与状态机库 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.tomlpip install pydantic-graph # 或使用 uv uv add pydantic-graph其运行环境要求与依赖如下项目要求Python 版本3.103.10 至 3.14 均受支持anyio4.7.0取消投递/流式清理依赖见 pyproject 注释logfire-api3.14.1OpenTelemetry 可观测性pydantic2.12typing-inspection0.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_typeint, output_typeint) 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(inputs4) 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()会直接抛出GraphSetupErrorNode ... is missing a return type hint on itsrunmethod。同时basenode.py 的文档字符串明确警告run方法的返回类型会在运行时被 pydantic-graph 读取用于定义接下来可以调用哪些节点并在图执行时强制执行。三、核心抽象BaseNode、End、GraphRunContext 与 Edgebasenode.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_onlyTrue) class GraphRunContext(Generic[StateT, DepsT]): state: StateT # 图的当前状态 deps: DepsT # 图运行依赖如数据库句柄、配置对象state在整个图运行期间共享是节点间传递可变数据的官方通道deps用于注入不随运行变化的外部依赖。End图的终止信号dataclass class End(Generic[RunEndT]): data: RunEndT # 图运行结束后对外返回的数据Edge边标签dataclass(frozenTrue) class Edge: label: str | None # 边的标签主要用于 Mermaid 渲染与可读性Edge是一个标注类型用来给图中的边打标签方便调试与可视化。四、GraphBuilder构建器的完整 API除了声明式BaseNodepydantic-graph 从 1.0 版本起主推构建器 API。GraphBuildergraph_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_instrumentTrue时每次run都会自动生成形如run graph name的 span每个步骤节点生成run node id的 span见 graph_builder.py 与 graph_builder.py。节点与边构建方法一览方法作用关键参数start_node/end_node获取图的入口/出口节点属性—step(callNone, *, node_idNone, labelNone)把异步函数包装成步骤节点可作装饰器或直接调用node_id缺省取函数名stream(callNone, *, node_idNone, labelNone)把返回异步迭代器的函数包装为流式步骤节点返回值类型为AsyncIterable[OutputT]node(node_type)将BaseNode子类接入图自动分析run返回注解建边缺失返回注解抛GraphSetupErrorjoin(reducer, *, initial/initial_factory, node_idNone, parent_fork_idNone, preferred_parent_forkfarthest)创建汇聚节点见并行与汇聚一节add(*edges)一次性加入多条边路径自动创建 fork/decision 等中间节点可传EdgePathadd_edge(source, destination, *, labelNone)添加一条简单边—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(*, noteNone, node_idNone)创建条件分支节点note会渲染进 Mermaid 图match(source, *, matchesNone)创建基于类型/谓词的分支匹配器见条件分支match_node(source, *, matchesNone)针对BaseNode子类的分支匹配—build(validate_graph_structureTrue)完成建图并返回可执行的Graph见构建校验README 示例中g.edge_from(g.start_node).to(start)便是edge_fromto的典型用法从__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_structureFalse关闭校验图中不存在死节点、结束节点必须可达、所有节点必须从 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)GraphRungraph_builder.py是单次图执行的实例暴露以下关键成员成员说明next(valueNone)推进一步执行返回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(initFalse) 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 / StreamFunctionStepFunction是异步函数协议async def f(ctx: StepContext) - OutputTStreamFunction是异步迭代器协议async def f(ctx: StepContext) - AsyncIterator[OutputT]配合GraphBuilder.stream()使用可在节点间流式产出数据Step封装函数 节点 ID 标签其as_node(inputsNone)返回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 与 matchdecision.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.pysource类型匹配方式Any/object恒为真该分支总是匹配Literal[...]输入值属于字面量集合之一即匹配其他类型用isinstance(inputs, source)判定分支按顺序测试命中第一个匹配分支的路径执行。若所有分支都不匹配运行时抛出RuntimeError: No branch matched inputs ...。分支的路径修饰DecisionBranchBuilder还提供链式方法.to(dest, *extra_dests, fork_idNone)指定一个或多个目标多个目标自动生成广播 fork.broadcast(get_forks)把分支广播到由回调产生的多个分支路径.transform(func)在分支路径上施加同步数据转换TransformFunction接收StepContext返回新值.map(fork_idNone, downstream_join_idNone)把可迭代输出摊平为多条并行路径.label(label)为路径段打标签仅用于 Mermaid 渲染。Decision还通过HandledT逆变类型参数在静态类型检查层面保证所有可能的输入类型都被穷尽处理_force_handled_contravariant方法decision.py是这一穷尽性检查的类型层面实现细节。八、并行与汇聚Fork、Map、Join 与 Reducer并行是 pydantic-graph 的杀手级能力其机制在 node.py 与 join.py 中定义。Fork广播 vs 映射Fork节点有两种模式node.pyis_map行为输入要求False广播 broadcast同一份数据被派发到每一条下游路径InputT OutputTTrue映射 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_idNone, preferred_parent_forkfarthest, # 或 closest )Join.reduce()会根据 reducer 的形参个数自动区分两种签名join.py普通 reducer(current, inputs) - current上下文 reducer(ctx: ReducerContext, current, inputs) - current内置 reducer 一览join.pyReducer行为reduce_null丢弃所有输入返回None仅作汇聚信号reduce_list_append追加单个元素到列表reduce_list_extend用可迭代扩展列表reduce_dict_update用映射更新字典reduce_sum数值累加ReduceFirstValue取第一个到达的值并取消其余兄弟任务早期停止ReducerContext暴露state、deps与cancel_sibling_tasks()——后者允许实现首个结果即返回的短路语义。ReduceFirstValue正是通过调用它实现取消的。支配 forkDominating ForkJoin 的同步前提每一个 Join 都必须满足存在一个支配它的 Fork F从 StartNode 到该 Join 的所有路径都必须经过 F且不含绕过 F 的环。构建器通过_collect_dominating_forksgraph_builder.py自动计算若找不到支配 fork 会抛出GraphBuildingError并附带 Mermaid 图帮助排查。这一约束保证引擎能判断该 fork 的所有上游任务是否已完成。preferred_parent_fork参数farthest/closest则在存在多个候选 fork 时决定优先选择哪一个。九、Mermaid 可视化一张图看懂你的图Graph.render()把整个图渲染为 Mermaid 的stateDiagram-v2字符串graph_builder.pystr(graph)也直接返回渲染结果因此调试时只需print(graph)。mermaid_text fives_graph.render( title逢 5 递进图, directionLR, # TB(默认) | LR | RL | BT )渲染细节graph_builder.pystart/end 节点用 Mermaid 的[*]特殊语法表示step 节点输出为node_id: labeljoin 节点输出为state id joinfork 输出为forkdecision 输出为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_structureTrue)默认执行结构校验graph_builder.py检查项包括start 节点必须存在出边图中必须存在通向 end 节点的边除 end 外不允许死端节点无出边、非决策分支end 节点必须从 start 可达所有节点必须从 start 可达。若你的图刻意违反上述假设例如某些节点允许成为死端可以显式传validate_graph_structureFalse跳过校验。校验失败时抛出的GraphValidationError会附带提示语If this is intentional, you can suppress this error by passingvalidate_graph_structureFalse...异常体系定义在 exceptions.py共五类异常基类触发场景GraphSetupErrorTypeError图配置错误如节点run缺失返回注解GraphBuildingErrorValueError建图阶段错误如节点 ID 冲突、Join 缺少支配 forkGraphValidationErrorValueError图结构校验失败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可视化等可作为理解行为边界的权威示例集。十二、快速上手清单安装pip install pydantic-graphPython ≥ 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),仅供参考
返回列表