pydantic-graph 图构建器 API 完全指南:GraphBuilder、Graph 与 GraphRun 实战详解
【免费下载链接】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-ai 仓库中pydantic_graph.graph_builder模块(即 docs/api/pydantic_graph/graph_builder.md 对应的 API 参考)展开,系统讲解基于构建器(Builder)的图工作流 API:如何用GraphBuilder以声明式方式构建可执行图,用Graph/GraphRun执行它,并通过 Mermaid 渲染图结构。读完本文,你将掌握类型化图工作流(状态、依赖、输入输出全类型标注)的完整构建、校验、执行、逐步迭代与可视化方法,能够把多步 Agent 流程、并行分支、决策路由与汇合聚合落地为可运行的图。
模块定位:构建式图 API 的“正统”入口
pydantic_graph.graph_builder是整个 pydantic-graph 中构建式图 API 的核心模块。模块 docstring 明确写道:它是构建式图 API 的 canonical home,即:
GraphBuilder:声明式构造可执行图;Graph与GraphRun:执行图;- Mermaid 渲染辅助函数:供
Graph.render()使用。
同一批公开符号同时从pydantic_graph顶层直接 re-export(见 pydantic_graph/pydantic_graph/init.py 中from .graph_builder import ...与__all__),因此你可以直接from pydantic_graph import GraphBuilder, Graph, GraphRun。整个 pydantic-ai 的 Agent 循环本身也是由 pydantic-graph 驱动(见 pydantic_graph/pydantic_graph/init.py 的说明),可见该模块的底层地位。
模块内部结构(graph_builder.py,共 2300 余行)大体分为四个部分:
- 图执行器:
EndMarker、ErrorMarker、Graph、GraphRun、GraphTask等运行时组件; - 图构建器:
GraphBuilder及其节点/边构建方法; - 构建期处理函数:占位 ID 替换、路径展平、fork 归一化、结构校验、支配 fork 收集等;
- Mermaid 渲染:
MermaidGraph、MermaidNode、MermaidEdge与拓扑排序。
对应的行为测试集中在 tests/graph/builder/ 目录,例如 test_graph_builder.py、test_joins_and_reducers.py、test_decisions.py 等,可作为每个 API 的“可运行示例库”。
GraphBuilder:声明式构建图的入口
GraphBuilder(graph_builder.py#L1138-L1149)是一个泛型类,类型参数定义了整张图的类型契约:
| 类型参数 | 含义 |
|---|---|
StateT | 图状态类型,贯穿整个图运行过程的可变状态 |
DepsT | 依赖类型,例如数据库连接、外部服务客户端 |
GraphInputT | 图输入数据(start_node接收的初始输入) |
GraphOutputT | 图输出数据(end_node产出的最终结果) |
构造函数的参数(graph_builder.py#L1183-L1217):
name: str | None = None:图的可选名称。若未提供,会在第一次调用图方法时从调用帧推断名称(infer_obj_name)。state_type/deps_type/input_type/output_type:均接受TypeOrTypeExpression(类型或类型表达式,如Literal[...]),默认值为NoneType。auto_instrument: bool = True:是否自动创建可观测性 span。开启后,Graph.run会创建名为run graph {name}的 span,每个Step节点执行时创建run node {node.id}的 span(graph_builder.py#L344-L352 与 graph_builder.py#L903-L908)。
GraphBuilder内部维护_nodes(节点字典)、_edges_by_source(按源节点索引的出边路径)、_decision_index,并预置了_start_node = StartNode()与_end_node = EndNode()。start_node/end_node属性暴露它们,其 ID 固定为__start__与__end__(见 node.py#L26-L43)。
最小可运行示例
以下示例来自 tests/graph/builder/test_graph_builder.py#L24-L42 的核心形态:
from dataclasses import dataclass from pydantic_graph import GraphBuilder, StepContext @dataclass class SimpleState: counter: int = 0 result: str | None = None g = GraphBuilder(state_type=SimpleState, output_type=int) @g.step async def increment(ctx: StepContext[SimpleState, None, None]) -> int: ctx.state.counter += 1 return ctx.state.counter g.add( g.edge_from(g.start_node).to(increment), g.edge_from(increment).to(g.end_node), ) graph = g.build() state = SimpleState() result = await graph.run(state=state) assert result == 1 assert state.counter == 1模式非常清晰:@g.step把异步函数包装成Step节点;g.edge_from(source).to(destination)构建边;g.add(...)一次性注册;g.build()生成可执行Graph;await graph.run(state=...)执行并返回最终输出。
节点构建:step / stream / join / decision / match
GraphBuilder提供多种节点构建方法,覆盖工作流中的常见控制流单元。
step:异步步骤节点
step(graph_builder.py#L1238-L1289)既可作为装饰器,也可直接调用:
# 装饰器形式 @g.step(node_id='custom_step_id', label='My Custom Label') async def my_step(ctx: StepContext[SimpleState, None, None]) -> int: return 42 # 直接调用形式 step = g.step(my_async_func, node_id='custom_step_id', label='My Custom Label')参数:node_id指定节点 ID(缺省时从函数名推断,get_callable_name);label是给人看的人类可读标签,会出现在 Mermaid 渲染中。步骤函数签名必须符合StepFunction协议——接收StepContext[StateT, DepsT, InputT],返回Awaitable[OutputT](step.py#L66-L87)。StepContext通过只读属性暴露state、deps、inputs(step.py#L25-L63)。
stream:异步迭代器流节点
stream(graph_builder.py#L1291-L1363)与step类似,但包装的是StreamFunction——一个返回AsyncIterator[OutputT]的异步可调用(step.py#L90-L112)。实现上,stream会把该调用包进一个async def wrapper(ctx)以统一执行路径,随后同样委托给self.step(...)。适合“逐步产出”的场景,例如逐块读取、流式生成。
join:汇合并行分支
join(graph_builder.py#L1365-L1405)创建Join节点,用 reducer 函数聚合来自多个并行分支的数据:
joined = g.join( reduce_list_append, initial=[], node_id='collect_results', parent_fork_id=None, preferred_parent_fork='farthest', )关键参数:
reducer:ReducerFunction[StateT, DepsT, InputT, OutputT],即(current, inputs) -> new_current或带上下文的(ctx, current, inputs) -> new_current两种形态(join.py#L81-L98)。Join.reduce会通过inspect.signature检查参数个数来区分两种形态(join.py#L194-L199)。initial或initial_factory:reducer 的初始值。两者互斥语义,缺省时initial_factory退化为lambda: initial。initial_factory适合每次运行需要新可变对象(如空 list/dict)的场景。node_id:节点 ID,缺省时基于 reducer 名生成占位 ID。parent_fork_id:显式指定父 fork;preferred_parent_fork:'farthest' | 'closest',当存在多个候选支配 fork 时选择最远或最近的(默认'farthest')。
内置 reducer 位于 join.py#L101-L147:
| reducer | 行为 |
|---|---|
reduce_null | 丢弃所有输入,返回None |
reduce_list_append | 把单个元素追加进 list |
reduce_list_extend | 把可迭代对象扩展进 list |
reduce_dict_update | 用 Mapping 更新 dict |
reduce_sum | 数值求和(要求类型支持__add__) |
ReduceFirstValue | 返回第一个到达的值,并取消其余兄弟任务(提前终止语义) |
ReduceFirstValue是@dataclass,可实例化后作为 reducer 传入:g.join(ReduceFirstValue(), initial=...)。它在内部调用ctx.cancel_sibling_tasks()触发提前终止(join.py#L140-L147)。任何 reducer 也都可以通过ReducerContext.cancel_sibling_tasks()自行实现“早停”,运行时_cancel_sibling_tasks会取消同 fork 下尚未完成的任务(graph_builder.py#L1096-L1106)。
decision / match:类型安全的条件路由
decision(graph_builder.py#L1531-L1541)创建一个空Decision节点,随后用.branch(...)追加分支;match(graph_builder.py#L1543-L1563)创建一个分支匹配器,基于source类型或自定义matches谓词路由:
decision = g.decision(note='route based on input type') decision = decision.branch( g.match(int).to(int_handler) # inputs 是 int 时走 int_handler ) decision = decision.branch( g.match(str).to(str_handler) # inputs 是 str 时走 str_handler ) g.add(g.edge_from(g.start_node).to(decision), ...)运行时_handle_decision(graph_builder.py#L925-L948)按分支顺序测试输入:若提供了matches谓词则直接调用;否则按source类型表达式判定——Any/object恒匹配、Literal[...]用成员判断、其他类型用isinstance。没有任何分支匹配时抛RuntimeError。DecisionBranchBuilder还支持链式.transform(...)(同步变换)、.map(...)(逐项并行展开)、.broadcast(...)(广播到多条路径)、.label(...)(仅供 Mermaid 渲染的标签),见 decision.py#L134-L276。
match_node / node:与声明式 BaseNode 无缝衔接
match_node(graph_builder.py#L1565-L1583):针对BaseNode子类做匹配,返回DecisionBranch[SourceNodeT],适合在旧式BaseNode.run返回类型上做分发。node(graph_builder.py#L1586-L1619):把一个BaseNode子类接入构建式图。它读取node_type.run的返回类型注解(get_type_hints),据此自动推断出边;若缺少返回类型注解则抛GraphSetupError。这一步依赖BaseNode.run的返回注解在运行时被读取并用于约束后续节点(basenode.py#L42-L44)。
此外,Step.as_node()返回StepNode、Join.as_node()返回JoinNode,让BaseNode子类可以把执行权交给构建式的Step/Join(见 step.py#L150-L198 与 join.py#L201-L248)。
边构建:edge_from、add 与路径标记
图的“边”在 pydantic-graph 中是一段Path(路径),由若干PathItem标记构成(paths.py#L60-L159):
| 标记 | 作用 |
|---|---|
TransformMarker | 在路径中同步变换数据(transform) |
MapMarker | 把可迭代输入逐项并行分发(map),即“散开” |
BroadcastMarker | 把同一份数据广播到多条并行路径(broadcast) |
LabelMarker | 为路径段加标签,仅用于 Mermaid 渲染(label) |
DestinationMarker | 路径终点,指向目标节点(to) |
edge_from:构建边的起点
edge_from(*sources)(graph_builder.py#L1518-L1529)返回EdgePathBuilder,支持链式调用:
g.add( g.edge_from(g.start_node).to(step_a), # 简单边 g.edge_from(g.start_node).to(step_b, step_c), # 多目标 => 自动广播 fork g.edge_from(step_a).transform(fmt).to(step_b), # 带变换 g.edge_from(step_a).map().to(step_b), # 逐项并行 g.edge_from(step_a).label('passing data').to(step_b), # 带标签 )EdgePathBuilder.map()有一个值得注意的限制:当前不支持多源节点上的 map,源码会直接抛NotImplementedError(paths.py#L387-L395),提示为每个源单独建边。
add:注册边并自动补全节点
add(*edges)(graph_builder.py#L1408-L1472)接受一个或多个EdgePath:
- 为每个
source插入节点、登记出边; - 递归处理
destination(含Decision分支内的嵌套目标,用destination_ids集合防环); - 路径中出现
BroadcastMarker/MapMarker时自动创建Fork节点并插入; - 自动边推断:对每个
Step目的地,用get_type_hints(destination.call, ...)读取返回注解,若返回类型可解析为节点/End/StepNode/JoinNode/BaseNode子类,则自动生成边(_edge_from_return_hint,见 graph_builder.py#L1641-L1708)。返回StepNode/JoinNode时必须用Annotated[...]携带对应Step/Join实例,否则抛GraphSetupError。
add_edge / add_mapping_edge:便捷封装
add_edge(source, destination, *, label=None)(graph_builder.py#L1474-L1485):单条简单边的快捷方式,等价于edge_from(source).label(...).to(destination)后再add。add_mapping_edge(source, map_to, *, pre_map_label=None, post_map_label=None, fork_id=None, downstream_join_id=None)(graph_builder.py#L1487-L1514):为“可迭代数据逐项并行处理”封装,支持 map 前后标签与显式 fork/join ID。downstream_join_id尤其重要:当映射的迭代器为空时,运行时仍会以初始值触发该 join(见_handle_fork_edges中的 eager 创建逻辑,graph_builder.py#L1044-L1055),避免空列表导致 join 永远不触发。
build():从描述到可执行图
build(validate_graph_structure: bool = True)(graph_builder.py#L1711-L1749)把累积的节点与边转换成一个可执行的Graph。构建期处理管线依次为:
_replace_placeholder_node_ids:把decision/match/map/broadcast生成的占位 ID 替换为稳定 ID(同名冲突时追加_2、_3后缀),见 graph_builder.py#L2062-L2089;_flatten_paths:在第一个MapMarker/BroadcastMarker处拆分路径,把并行段从路径中抽离为独立边,见 graph_builder.py#L1862-L1901;_normalize_forks:归一化图结构,保证只有广播 fork 才有多条出边——任何有多条出边的普通节点都会被自动插入一个{node.id}_broadcast_fork,见 graph_builder.py#L1904-L1939;_validate_graph_structure:结构校验(见下);_collect_dominating_forks:为每个Join计算支配 fork,见 graph_builder.py#L1942-L2025;_compute_intermediate_join_nodes:计算 join 之间的“中间 join”关系,用于判定哪些 join 是“最终的”(final),见 graph_builder.py#L2028-L2059。
图结构校验规则
_validate_graph_structure(graph_builder.py#L1752-L1859)会检查五类问题,任一不满足即抛GraphValidationError:
- start 节点必须存在出边;
- 必须存在到达 end 节点的边;
- 除 end 节点外,不允许存在“死胡同”节点(无出边的非 end 节点);
- end 节点必须从 start 节点可达;
- 所有节点必须从 start 节点可达。
若确实需要构建违反上述假设的图(例如故意留死路),可传validate_graph_structure=False关闭校验——错误信息中会附上这句提示。
支配 fork(dominating fork)约束
构建期会为每个Join求“支配 fork”:从 start 到该 join 的所有路径都必须经过它、且包含该 join 的环也必须经过它。求不出来时抛GraphBuildingError,错误信息中会附上整张图的 Mermaid 渲染结果,方便定位。这是 join 语义得以成立的前提:运行时靠它判断“该 fork 上游的所有任务是否都已完成,可以继续向下游执行”。
Graph 与 GraphRun:图的执行模型
Graph:一次构建、多次执行
Graph(graph_builder.py#L157-L387)是完整可执行的图定义,字段包括name、state_type、deps_type、input_type、output_type、auto_instrument、nodes(按 ID 索引的节点字典)、edges_by_source、parent_forks(每个 join 的父 fork 信息)与intermediate_join_nodes。其辅助方法get_parent_fork(join_id)与is_final_join(join_id)(graph_builder.py#L205-L238)在运行期被频繁调用。
执行入口有三个:
| 方法 | 语义 | 使用限制 |
|---|---|---|
await graph.run(state=..., deps=..., inputs=...) | 异步执行到结束,返回最终输出(graph_builder.py#L240-L279) | 内部用graph_run.next(...)循环驱动,直到收到EndMarker |
graph.run_sync(...) | 同步便捷封装,基于loop.run_until_complete(graph_builder.py#L281-L310) | 不能在异步代码或已有运行中事件循环的环境里调用 |
graph.iter(...) | 异步上下文管理器,产出GraphRun供逐步执行(graph_builder.py#L312-L361) | 需要细粒度控制时使用 |
三个方法的参数一致:state、deps、inputs均为关键字参数,另有span(外部 span 上下文)与infer_name(是否从调用帧推断图名)。未提供span且auto_instrument=True时,iter会进入logfire_span(f'run graph {self.name}', graph=self),并把traceparent传给GraphRun用于链路追踪。名称推断深度在不同入口有差异(run用 depth=2,iter因asynccontextmanager包装用 depth=3)。
Graph.__str__返回 Mermaid 图文本,__repr__则在默认表示中嵌入__str__结果,方便 REPL 调试。
GraphRun:单次执行的运行时
GraphRun(graph_builder.py#L430-L634)管理一次执行的完整运行时状态:任务调度、fork/join 协调、结果追踪。核心成员:
state/deps/inputs:本次运行的上下文;__aiter__/__anext__:原生异步迭代,每次产出EndMarker[OutputT] | Sequence[GraphTask];next(value=None):推进一个步骤(graph_builder.py#L564-L584),可传入新的任务序列或EndMarker覆盖下一步;override_next(value):在End或节点报错后重定向执行(graph_builder.py#L586-L598),被after_node_run、on_node_run_error等钩子系统用于注入新任务或提前结束;只能在两次迭代之间调用;next_task属性:下一批待执行任务(未设置时返回首个任务);output属性:若图已完成,返回最终输出,否则为None。
错误处理上,节点抛出的异常会以ErrorMarker形式通过迭代器产出(而不是直接抛出),调用方可在下一次迭代前用override_next恢复;若调用方不做处理,__anext__会把错误重新抛出。Graph.run的内部循环会在StopAsyncIteration时断言最后一次事件必须是EndMarker,否则视为运行器 bug。
运行时的核心调度在_GraphIterator.iter_graph(graph_builder.py#L673-L841):通过 anyio 任务组并发执行GraphTask,结果经内存流(create_memory_object_stream)回传;join 项目(JoinItem)按 fork 归属归入active_reducers,由_resolve_join_fork_run确定归属的 fork run;当某 fork run 的所有任务完成(_is_fork_run_completed)且没有更深的中间 join 等待时,才 finalize 该 join 并把结果沿下游边继续派发。整体上是一套“fork 并发 → join 归约 → 下游继续”的同步屏障模型。
并行控制流:Fork 与 Join 的协作机制
并行能力由Fork节点提供(node.py#L60-L79)。Fork有两种模式:
is_map=False(广播):InputT即OutputT,同一份数据发给所有分支;is_map=True(映射):InputT必须是Sequence[OutputT]或异步可迭代对象,每个元素进入一条独立分支。
运行时_handle_fork_edges(graph_builder.py#L1033-L1082)为每个分支生成GraphTask,并在 fork 栈(ForkStack)中记录(fork_id, node_run_id, thread_index)元组,这是后续 join 判定“哪些任务属于同一 fork run”的依据。映射支持同步可迭代与异步可迭代两种输入;输入既不可迭代也不可异步迭代时抛出RuntimeError('Cannot map non-iterable value: ...')。
Join的同步语义(join.py#L150-L199):同一 fork run 的多个JoinItem依次喂给 reducer,reducer 的中间结果保存在JoinState中;只有该 fork run 的所有上游任务都完成(_is_fork_run_completed,graph_builder.py#L1084-L1094)时,join 才被 finalize 并继续下游。对于嵌套 join(一个 join 位于另一个 join 与父 fork 之间),_compute_intermediate_join_nodes与is_final_join共同决定:
- final join:合并时截断 fork 栈到父 fork 为止(
_resolve_join_fork_run,graph_builder.py#L967-L981); - 非 final join:保留完整 fork 栈,使下游 join 仍能与同一 fork run 关联。
无活动任务时的收尾阶段,迭代器会按“先 finalize 无中间 join 的 reducer”的顺序推进(graph_builder.py#L765-L834),避免中间 join 的输出还没到达就先处理了外层 join。
可视化:Graph.render() 与 Mermaid 渲染管线
Graph.render(title=None, direction=None)(graph_builder.py#L363-L373)返回 Mermaid 状态图字符串,实际由模块级函数build_mermaid_graph(nodes, edges_by_source).render(...)完成。
节点与边的映射
build_mermaid_graph(graph_builder.py#L2168-L2219)把图中每种节点映射为 Mermaid 语法:
| 节点类型 | Mermaid 呈现 |
|---|---|
StartNode/EndNode | 边两端用[*]特殊语法 |
Step | {id}: {label} |
Join | state {id} <<join>> |
Fork(map / broadcast) | state {id} <<fork>> |
Decision | state {id} <<choice>>,note通过note right of呈现 |
边上的LabelMarker标签会输出为A --> B: label的注释文本。
渲染参数与方向
StateDiagramDirection(graph_builder.py#L2137-L2144)定义了四种布局方向:
'TB':自上而下(Mermaid 默认);'LR':自左而右;'RL':自右而左;'BT':自下而上。
MermaidGraph.render(graph_builder.py#L2232-L2279)支持title(输出---\ntitle: ...\n---前置块)、direction与edge_labels(默认 True)三个参数;输出以stateDiagram-v2开头。渲染前会做拓扑排序(_topological_sort,graph_builder.py#L2282-L2324):从 start 节点做 BFS 求各节点深度,按“距 start 的距离”排序节点与边,保证图中元素呈现自然的“上游在前、下游在后”顺序。
用法示例:
graph = g.build() print(graph) # 等价于 print(graph.render()) print(graph.render(title='my graph', direction='LR'))此外,构建期的GraphBuildingError(如 join 缺少支配 fork)也会附带渲染好的 Mermaid 图,帮助快速定位结构问题(graph_builder.py#L2008-L2022)。
类型系统与可观测性的设计要点
从源码可以提炼出该模块三个贯穿始终的设计原则:
- 类型即契约:四个泛型参数贯穿
GraphBuilder → Graph → GraphRun全程;Decision通过HandledT的逆变约束做分支穷尽性静态检查(decision.py#L68-L80);BaseNode.run的返回注解在运行期被读取用于自动建边。这意味着静态类型检查器(pyright / mypy)能帮你提前发现“某输入类型没有对应分支”“返回了未声明的下游节点”等错误。 - 构建期重写,运行期简化:占位 ID 替换、路径展平、fork 归一化都在
build()期间完成,运行期_GraphIterator面对的是一张已被规范化的图(例如所有多出边都被折叠进显式广播 fork),从而显著简化执行逻辑。 - 可观测性内置:
auto_instrument默认开启,run graph {name}与run node {id}两级 span 覆盖整次运行与单个步骤;traceparent沿GraphRun传递,保证与外部 tracing 系统的链路一致性。模块还通过_unwrap_exception_groups把异常组解包为原始异常,避免日志与错误处理被ExceptionGroup包裹(graph_builder.py#L1117-L1132)。
从示例到实战:推荐的学习路径
如果你想把本文内容落到手头代码里,建议按以下顺序在仓库中对照研读:
- 最小链路:test_graph_builder.py(顺序步骤、输入传递、自定义 ID 与 label);
- 并行与汇合:test_joins_and_reducers.py、test_broadcast_and_spread.py(map/broadcast、各种 reducer、空迭代器与提前终止);
- 条件路由:test_decisions.py(decision/match/transform/map 组合);
- 逐步执行与钩子:test_graph_iteration.py(
GraphRun.next/override_next/ErrorMarker恢复); - 旧式节点融合:test_basenode_integration.py(
node()、as_node()、StepNode/JoinNode桥接); - 边界情况:test_edge_cases.py、test_graph_edge_cases.py(重复节点 ID、无匹配分支、多分支 return hint 等)。
图的基础概念(节点、边、状态、依赖)与更多使用场景,可进一步参考 docs/graph/builder/index.md 与 docs/graph.md;构建器相关 API 的类型注解与异常定义,可在 pydantic_graph/pydantic_graph/init.py 的__all__中一览全貌,而GraphSetupError/GraphRuntimeError/GraphBuildingError/GraphValidationError的具体语义定义在 pydantic_graph/pydantic_graph/exceptions.py。
掌握GraphBuilder的构建、Graph的执行与GraphRun的细粒度控制之后,你便拥有了一套类型安全、可并行、可观测、可直接可视化的图工作流引擎,可以把它作为多步 Agent、数据处理流水线乃至任何有向工作流的基础设施来使用。
【免费下载链接】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),仅供参考