[LangGraph编译原理]从“状态图(StateGraph)”到“Actor模型(Pregel)”的转变
我们通过添加节点和边创建StateGraph对象其compile方法将其编译生成CompiledStateGraph整个编译过程就是将StateGraph的节点和边转换成Pregel的节点和通道的过程。整个编译流程主要划分为三个步骤。步骤一 创建CompiledStateGraph第一个步骤就是按照如下的方式将CompiledStateGraph对象创建出来。classStateGraph(Generic[StateT,ContextT,InputT,OutputT]):defcompile(self,checkpointer:CheckpointerNone,*,cache:BaseCache|NoneNone,store:BaseStore|NoneNone,interrupt_before:All|list[str]|NoneNone,interrupt_after:All|list[str]|NoneNone,debug:boolFalse,name:str|NoneNone,)-CompiledStateGraph[StateT,ContextT,InputT,OutputT]:checkpointerensure_valid_checkpointer(checkpointer)# assign default valuesinterrupt_beforeinterrupt_beforeor[]interrupt_afterinterrupt_afteror[]# validate the graphself.validate(interrupt((interrupt_beforeifinterrupt_before!*else[])interrupt_afterifinterrupt_after!*else[]))# prepare output channelsoutput_channels(__root__iflen(self.schemas[self.output_schema])1and__root__inself.schemas[self.output_schema]else[keyforkey,valinself.schemas[self.output_schema].items()ifnotis_managed_value(val)])stream_channels(__root__iflen(self.channels)1and__root__inself.channelselse[keyforkey,valinself.channels.items()ifnotis_managed_value(val)])compiledCompiledStateGraph[StateT,ContextT,InputT,OutputT](builderself,schema_to_mapper{},context_schemaself.context_schema,nodes{},channels{**self.channels,**self.managed,START:EphemeralValue(self.input_schema),},input_channelsSTART,stream_modeupdates,output_channelsoutput_channels,stream_channelsstream_channels,checkpointercheckpointer,interrupt_before_nodesinterrupt_before,interrupt_after_nodesinterrupt_after,auto_validateFalse,debugdebug,storestore,cachecache,namenameorLangGraph,)...创建CompiledStateGraph对象调用构造函数的参数来源于StateGraph自身相关的字段和调用compile方法传入的参数。它会对如下几个输入作相应的处理表示节点列表的nodes将在下一步骤添加此时暂时设置为空表示通道和ManagedValue的channels是通过解析状态Schema类型生成。在这些针对状态成员的通道列表上还会外加上一个名为__start__的通道类型为EphemeralValue。__start__被设置为输入通道意味着我们调用Agent提供的输入整个会被写入此通道输出通道通过解析输出Schema类型生成流式输出通道为channels中剔除ManagedValue余下的部分步骤二附加节点第一步创建的CompiledStateGraph对象其节点为空这部分通过第二个步骤来补充。如下面的代码片段所示在完成CompiledStateGraph的创建之后compile方法除了会调用后者的attach_node方法附加所有节点之外还会额外附加一个名为__start__的节点。也就是说我们在构建图的时候作为开始和终结的节点__start__和__end__只有前者才对应一个真正的节点PregelNode后者可以视为一个代表“虚拟终结节点”。之所以需要添加一个额外的__start__节点是因为作为Actor模型的Pregel需要一个统一的“启动点”类似我们熟悉的Main函数。试想一下如果没有这个其实节点我们调用Agent时必需知道当前Pregel的入口节点是哪些并需要针对它们的输入通道提供对应的输入。如果Agent有明确的静态入口节点还好但是很多情况下这个起始节点是动态根据状态决定的所以我们才需要set_conditional_entry_point方法处理起来就很麻烦了。classStateGraph(Generic[StateT,ContextT,InputT,OutputT]):defcompile(self,checkpointer:CheckpointerNone,*,cache:BaseCache|NoneNone,store:BaseStore|NoneNone,interrupt_before:All|list[str]|NoneNone,interrupt_after:All|list[str]|NoneNone,debug:boolFalse,name:str|NoneNone,)-CompiledStateGraph[StateT,ContextT,InputT,OutputT]:#步骤一 创建CompiledStateGraphcompiled.attach_node(START,None)forkey,nodeinself.nodes.items():compiled.attach_node(key,node)...attach_node方法创建PregelNode的逻辑根据节点类型__start__节点和常规节点有所不同。为__start__节点创建PregelNode__start__节点会订阅同名__start__通道由于这是Pregel的输入通道所以我们调用Pregel时提供的整个输入都将进入此通道。节点需要作的仅仅是拆解输入将其分发给对应的通道。所以节点对应的PregelNode对象并不绑定任何操作其bound字段无需指定其核心部分仅仅是writes字段体现的针对输出通道的写入意图。最终被添加到writes字段的是一个根据ChannelWriteTupleEntry创建的ChannelWrite对象这个ChannelWriteTupleEntry利用其mapper将输入转换成一系列channel_name, channel_value二元组完成输入通道的分发。为常规节点创建PregelNode常规节点在StateGraph中被描述为StateNodeSpec对象常见的PregelNode正是根据StateNodeSpec创建而成具体每个字段的初始化规则如下bund该字段被设置为StateNodeSpec的runnable字段所以执行PregelNode就是执行对应的节点函数triggers每个节点具有各自独立的驱动通道并统一按照branch:to:{node_name}这样的格式命名所以该字段就是包含对应通道名称的单元素列表channels: 通过解析输入Schema类型得到mapper如果输入Schema类型和状态Schema不同那么将根据输入类型从schema_to_mapper字段中提取用于映射输入对象的Runnable作为该字段的值否则就是对等映射metadata/retry_policy/cache_policyStateNodeSpec对象的同名字段writes包含一个根据ChannelWriteTupleEntry对象创建的ChannelWrite对象ChannelWriteTupleEntry利用mapper字段按照如下的方式将针对状态的增量更新和节点跳转转换成一系列channel_name, channel_value二元组如果返回字典或者Pydantic模型针对其数据成员生成channel_name, channel_value二元组如果返回一个Command对象将其update字段按照上面的方式转换成channel_name, channel_value二元组如果goto字段是一个或者多个Send转换成“__pregel_tasks”,Send二元组列表如果goto字段表示的跳转的目标节点则转换成一组“branch:to:{node_name}”,None二元组。这些二元组合并在一起作为mapper映射函数的返回值作为节点的驱动通道也是在attach_node方法中创建并注册的。通道按照branch:to:{node_name}格式进行命名因为具有“邻步有效”的特性默认类型为EphemeralValue。如果添加节点时将defer参数设置为True此时会创建一个LastValueAfterFinish类型的通道实现延迟执行的目的。步骤三附加边目前为止根据StateNodeSpec创建的PregelNode依旧是不完整的它的writes注册的ChannelWrite仅仅通过提交的通道写入意图完成了对状态的增量修改以及利用Command的goto字段实现的跳转功能但是“边”体现的节点之间的依赖关系还未通过通道写入建立起来。这是编译流程最后一步需要完成的任务。classStateGraph(Generic[StateT,ContextT,InputT,OutputT]):defcompile(self,checkpointer:CheckpointerNone,*,cache:BaseCache|NoneNone,store:BaseStore|NoneNone,interrupt_before:All|list[str]|NoneNone,interrupt_after:All|list[str]|NoneNone,debug:boolFalse,name:str|NoneNone,)-CompiledStateGraph[StateT,ContextT,InputT,OutputT]:#步骤二附加节点forstart,endinself.edges:compiled.attach_edge(start,end)forstarts,endinself.waiting_edges:compiled.attach_edge(starts,end)forstart,branchesinself.branches.items():forname,branchinbranches.items():compiled.attach_branch(start,name,branch)如上面的代码片段所示对于edges字段保存的“一对一静态边”和waiting_edges字段保存的“多对一静态边”都会调用如下这个attach_edge方法。对于前者只需要更新下一节点对应的驱动通道就可以了对于后者要求后续节点在前序节点全部执行之后才开始执行LangGraph定义了NamedBarrierValueAfterFinish这种专门的通道类型来解决这种场景我的文章“Channel——驱动Node执行的原力”对所有的通道类型有过详细介绍。classCompiledStateGraph(Pregel[StateT,ContextT,InputT,OutputT],Generic[StateT,ContextT,InputT,OutputT],):defattach_edge(self,starts:str|Sequence[str],end:str)-None:ifisinstance(starts,str):# subscribe to start channelifend!END:self.nodes[starts].writers.append(ChannelWrite((ChannelWriteEntry(_CHANNEL_BRANCH_TO.format(end),None),)))elifend!END:channel_namefjoin:{.join(starts)}:{end}# register channelifself.builder.nodes[end].defer:self.channels[channel_name]NamedBarrierValueAfterFinish(str,set(starts))else:self.channels[channel_name]NamedBarrierValue(str,set(starts))# subscribe to channelself.nodes[end].triggers.append(channel_name)# publish to channelforstartinstarts:self.nodes[start].writers.append(ChannelWrite((ChannelWriteEntry(channel_name,start),)))在完成了针对NamedBarrierValueAfterFinish通道的创建之后通道名称格式为fjoin:{.join(starts)}:{end}会被添加后续节点的triggers字段表示的订阅通道中。前序节点对应PregelNode的writes中会添加相应的ChannelWrite提交针对该通道的写入意图。对于表示条件边的每个BranchSpec对象同样需要转换成针对路由目标节点驱动通道的写入。通过前面针对BranchSpec的介绍我们知道了它的run方法返回的Runnable对象会为我们完成所有的工作所以我们只需要将这个Runnable对象添加到作为代表源节点PregelNode的writes字段中就可以了。compile方法会将每个BranchSpec对象作为参数调用attach_branch方法该方法就是这么做的。至此表示Agent的CompiledStateGraph完成初始化整个编译过程结束。