pydantic-graph 图构建器 API 完全指南:GraphBuilder、Graph 与 GraphRun 实战详解
2026/9/13 12:20:11 网站建设 项目流程

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:声明式构造可执行图;
  • GraphGraphRun:执行图;
  • 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 余行)大体分为四个部分:

  1. 图执行器EndMarkerErrorMarkerGraphGraphRunGraphTask等运行时组件;
  2. 图构建器GraphBuilder及其节点/边构建方法;
  3. 构建期处理函数:占位 ID 替换、路径展平、fork 归一化、结构校验、支配 fork 收集等;
  4. Mermaid 渲染MermaidGraphMermaidNodeMermaidEdge与拓扑排序。

对应的行为测试集中在 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()生成可执行Graphawait 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通过只读属性暴露statedepsinputs(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', )

关键参数:

  • reducerReducerFunction[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)。
  • initialinitial_factory:reducer 的初始值。两者互斥语义,缺省时initial_factory退化为lambda: initialinitial_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没有任何分支匹配时抛RuntimeErrorDecisionBranchBuilder还支持链式.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()返回StepNodeJoin.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。构建期处理管线依次为:

  1. _replace_placeholder_node_ids:把decision/match/map/broadcast生成的占位 ID 替换为稳定 ID(同名冲突时追加_2_3后缀),见 graph_builder.py#L2062-L2089;
  2. _flatten_paths:在第一个MapMarker/BroadcastMarker处拆分路径,把并行段从路径中抽离为独立边,见 graph_builder.py#L1862-L1901;
  3. _normalize_forks:归一化图结构,保证只有广播 fork 才有多条出边——任何有多条出边的普通节点都会被自动插入一个{node.id}_broadcast_fork,见 graph_builder.py#L1904-L1939;
  4. _validate_graph_structure:结构校验(见下);
  5. _collect_dominating_forks:为每个Join计算支配 fork,见 graph_builder.py#L1942-L2025;
  6. _compute_intermediate_join_nodes:计算 join 之间的“中间 join”关系,用于判定哪些 join 是“最终的”(final),见 graph_builder.py#L2028-L2059。

图结构校验规则

_validate_graph_structure(graph_builder.py#L1752-L1859)会检查五类问题,任一不满足即抛GraphValidationError

  1. start 节点必须存在出边;
  2. 必须存在到达 end 节点的边;
  3. 除 end 节点外,不允许存在“死胡同”节点(无出边的非 end 节点);
  4. end 节点必须从 start 节点可达;
  5. 所有节点必须从 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)是完整可执行的图定义,字段包括namestate_typedeps_typeinput_typeoutput_typeauto_instrumentnodes(按 ID 索引的节点字典)、edges_by_sourceparent_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)需要细粒度控制时使用

三个方法的参数一致:statedepsinputs均为关键字参数,另有span(外部 span 上下文)与infer_name(是否从调用帧推断图名)。未提供spanauto_instrument=True时,iter会进入logfire_span(f'run graph {self.name}', graph=self),并把traceparent传给GraphRun用于链路追踪。名称推断深度在不同入口有差异(run用 depth=2,iterasynccontextmanager包装用 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_runon_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(广播):InputTOutputT,同一份数据发给所有分支;
  • 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_nodesis_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}
Joinstate {id} <<join>>
Fork(map / broadcast)state {id} <<fork>>
Decisionstate {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---前置块)、directionedge_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)。

类型系统与可观测性的设计要点

从源码可以提炼出该模块三个贯穿始终的设计原则:

  1. 类型即契约:四个泛型参数贯穿GraphBuilder → Graph → GraphRun全程;Decision通过HandledT的逆变约束做分支穷尽性静态检查(decision.py#L68-L80);BaseNode.run的返回注解在运行期被读取用于自动建边。这意味着静态类型检查器(pyright / mypy)能帮你提前发现“某输入类型没有对应分支”“返回了未声明的下游节点”等错误。
  2. 构建期重写,运行期简化:占位 ID 替换、路径展平、fork 归一化都在build()期间完成,运行期_GraphIterator面对的是一张已被规范化的图(例如所有多出边都被折叠进显式广播 fork),从而显著简化执行逻辑。
  3. 可观测性内置auto_instrument默认开启,run graph {name}run node {id}两级 span 覆盖整次运行与单个步骤;traceparent沿GraphRun传递,保证与外部 tracing 系统的链路一致性。模块还通过_unwrap_exception_groups把异常组解包为原始异常,避免日志与错误处理被ExceptionGroup包裹(graph_builder.py#L1117-L1132)。

从示例到实战:推荐的学习路径

如果你想把本文内容落到手头代码里,建议按以下顺序在仓库中对照研读:

  1. 最小链路:test_graph_builder.py(顺序步骤、输入传递、自定义 ID 与 label);
  2. 并行与汇合:test_joins_and_reducers.py、test_broadcast_and_spread.py(map/broadcast、各种 reducer、空迭代器与提前终止);
  3. 条件路由:test_decisions.py(decision/match/transform/map 组合);
  4. 逐步执行与钩子:test_graph_iteration.py(GraphRun.next/override_next/ErrorMarker恢复);
  5. 旧式节点融合:test_basenode_integration.py(node()as_node()StepNode/JoinNode桥接);
  6. 边界情况: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),仅供参考

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

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

立即咨询