pydantic-graph 图构建器 API 完全指南:GraphBuilder、Graph 与 GraphRun 实战详解
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声明式构建图的入口GraphBuildergraph_builder.py#L1138-L1149是一个泛型类类型参数定义了整张图的类型契约类型参数含义StateT图状态类型贯穿整个图运行过程的可变状态DepsT依赖类型例如数据库连接、外部服务客户端GraphInputT图输入数据start_node接收的初始输入GraphOutputT图输出数据end_node产出的最终结果构造函数的参数graph_builder.py#L1183-L1217name: 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}的 spangraph_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_typeSimpleState, output_typeint) 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(statestate) 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 / matchGraphBuilder提供多种节点构建方法覆盖工作流中的常见控制流单元。step异步步骤节点stepgraph_builder.py#L1238-L1289既可作为装饰器也可直接调用# 装饰器形式 g.step(node_idcustom_step_id, labelMy Custom Label) async def my_step(ctx: StepContext[SimpleState, None, None]) - int: return 42 # 直接调用形式 step g.step(my_async_func, node_idcustom_step_id, labelMy Custom Label)参数node_id指定节点 ID缺省时从函数名推断get_callable_namelabel是给人看的人类可读标签会出现在 Mermaid 渲染中。步骤函数签名必须符合StepFunction协议——接收StepContext[StateT, DepsT, InputT]返回Awaitable[OutputT]step.py#L66-L87。StepContext通过只读属性暴露state、deps、inputsstep.py#L25-L63。stream异步迭代器流节点streamgraph_builder.py#L1291-L1363与step类似但包装的是StreamFunction——一个返回AsyncIterator[OutputT]的异步可调用step.py#L90-L112。实现上stream会把该调用包进一个async def wrapper(ctx)以统一执行路径随后同样委托给self.step(...)。适合“逐步产出”的场景例如逐块读取、流式生成。join汇合并行分支joingraph_builder.py#L1365-L1405创建Join节点用 reducer 函数聚合来自多个并行分支的数据joined g.join( reduce_list_append, initial[], node_idcollect_results, parent_fork_idNone, preferred_parent_forkfarthest, )关键参数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。initial或initial_factoryreducer 的初始值。两者互斥语义缺省时initial_factory退化为lambda: initial。initial_factory适合每次运行需要新可变对象如空 list/dict的场景。node_id节点 ID缺省时基于 reducer 名生成占位 ID。parent_fork_id显式指定父 forkpreferred_parent_forkfarthest | closest当存在多个候选支配 fork 时选择最远或最近的默认farthest。内置 reducer 位于 join.py#L101-L147reducer行为reduce_null丢弃所有输入返回Nonereduce_list_append把单个元素追加进 listreduce_list_extend把可迭代对象扩展进 listreduce_dict_update用 Mapping 更新 dictreduce_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类型安全的条件路由decisiongraph_builder.py#L1531-L1541创建一个空Decision节点随后用.branch(...)追加分支matchgraph_builder.py#L1543-L1563创建一个分支匹配器基于source类型或自定义matches谓词路由decision g.decision(noteroute 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_decisiongraph_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_nodegraph_builder.py#L1565-L1583针对BaseNode子类做匹配返回DecisionBranch[SourceNodeT]适合在旧式BaseNode.run返回类型上做分发。nodegraph_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在路径中同步变换数据transformMapMarker把可迭代输入逐项并行分发map即“散开”BroadcastMarker把同一份数据广播到多条并行路径broadcastLabelMarker为路径段加标签仅用于 Mermaid 渲染labelDestinationMarker路径终点指向目标节点toedge_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源码会直接抛NotImplementedErrorpaths.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, *, labelNone)graph_builder.py#L1474-L1485单条简单边的快捷方式等价于edge_from(source).label(...).to(destination)后再add。add_mapping_edge(source, map_to, *, pre_map_labelNone, post_map_labelNone, fork_idNone, downstream_join_idNone)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_structuregraph_builder.py#L1752-L1859会检查五类问题任一不满足即抛GraphValidationErrorstart 节点必须存在出边必须存在到达 end 节点的边除 end 节点外不允许存在“死胡同”节点无出边的非 end 节点end 节点必须从 start 节点可达所有节点必须从 start 节点可达。若确实需要构建违反上述假设的图例如故意留死路可传validate_graph_structureFalse关闭校验——错误信息中会附上这句提示。支配 forkdominating fork约束构建期会为每个Join求“支配 fork”从 start 到该 join 的所有路径都必须经过它、且包含该 join 的环也必须经过它。求不出来时抛GraphBuildingError错误信息中会附上整张图的 Mermaid 渲染结果方便定位。这是 join 语义得以成立的前提运行时靠它判断“该 fork 上游的所有任务是否都已完成可以继续向下游执行”。Graph 与 GraphRun图的执行模型Graph一次构建、多次执行Graphgraph_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(...)循环驱动直到收到EndMarkergraph.run_sync(...)同步便捷封装基于loop.run_until_completegraph_builder.py#L281-L310不能在异步代码或已有运行中事件循环的环境里调用graph.iter(...)异步上下文管理器产出GraphRun供逐步执行graph_builder.py#L312-L361需要细粒度控制时使用三个方法的参数一致state、deps、inputs均为关键字参数另有span外部 span 上下文与infer_name是否从调用帧推断图名。未提供span且auto_instrumentTrue时iter会进入logfire_span(frun graph {self.name}, graphself)并把traceparent传给GraphRun用于链路追踪。名称推断深度在不同入口有差异run用 depth2iter因asynccontextmanager包装用 depth3。Graph.__str__返回 Mermaid 图文本__repr__则在默认表示中嵌入__str__结果方便 REPL 调试。GraphRun单次执行的运行时GraphRungraph_builder.py#L430-L634管理一次执行的完整运行时状态任务调度、fork/join 协调、结果追踪。核心成员state/deps/inputs本次运行的上下文__aiter__/__anext__原生异步迭代每次产出EndMarker[OutputT] | Sequence[GraphTask]next(valueNone)推进一个步骤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_graphgraph_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_mapFalse广播InputT即OutputT同一份数据发给所有分支is_mapTrue映射InputT必须是Sequence[OutputT]或异步可迭代对象每个元素进入一条独立分支。运行时_handle_fork_edgesgraph_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依次喂给 reducerreducer 的中间结果保存在JoinState中只有该 fork run 的所有上游任务都完成_is_fork_run_completedgraph_builder.py#L1084-L1094时join 才被 finalize 并继续下游。对于嵌套 join一个 join 位于另一个 join 与父 fork 之间_compute_intermediate_join_nodes与is_final_join共同决定final join合并时截断 fork 栈到父 fork 为止_resolve_join_fork_rungraph_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(titleNone, directionNone)graph_builder.py#L363-L373返回 Mermaid 状态图字符串实际由模块级函数build_mermaid_graph(nodes, edges_by_source).render(...)完成。节点与边的映射build_mermaid_graphgraph_builder.py#L2168-L2219把图中每种节点映射为 Mermaid 语法节点类型Mermaid 呈现StartNode/EndNode边两端用[*]特殊语法Step{id}: {label}Joinstate {id} joinForkmap / broadcaststate {id} forkDecisionstate {id} choicenote通过note right of呈现边上的LabelMarker标签会输出为A -- B: label的注释文本。渲染参数与方向StateDiagramDirectiongraph_builder.py#L2137-L2144定义了四种布局方向TB自上而下Mermaid 默认LR自左而右RL自右而左BT自下而上。MermaidGraph.rendergraph_builder.py#L2232-L2279支持title输出---\ntitle: ...\n---前置块、direction与edge_labels默认 True三个参数输出以stateDiagram-v2开头。渲染前会做拓扑排序_topological_sortgraph_builder.py#L2282-L2324从 start 节点做 BFS 求各节点深度按“距 start 的距离”排序节点与边保证图中元素呈现自然的“上游在前、下游在后”顺序。用法示例graph g.build() print(graph) # 等价于 print(graph.render()) print(graph.render(titlemy graph, directionLR))此外构建期的GraphBuildingError如 join 缺少支配 fork也会附带渲染好的 Mermaid 图帮助快速定位结构问题graph_builder.py#L2008-L2022。类型系统与可观测性的设计要点从源码可以提炼出该模块三个贯穿始终的设计原则类型即契约四个泛型参数贯穿GraphBuilder → Graph → GraphRun全程Decision通过HandledT的逆变约束做分支穷尽性静态检查decision.py#L68-L80BaseNode.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.pymap/broadcast、各种 reducer、空迭代器与提前终止条件路由test_decisions.pydecision/match/transform/map 组合逐步执行与钩子test_graph_iteration.pyGraphRun.next/override_next/ErrorMarker恢复旧式节点融合test_basenode_integration.pynode()、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),仅供参考