ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

[拆解LangChain执行引擎-13]回到过去,开启平行世界[上篇]

[拆解LangChain执行引擎-13]回到过去,开启平行世界[上篇] Pregel提供了update_state/aupdate_state和bulk_update_state/abulk_update_state方法以增量的方式修改部分状态。有一点需要明确的是这些方法并非直接修改某个持久化的Checkpoint而是在此基础上创建一个子Checkpoint所以历史是不会被篡改的只会在某个时间点开启了一个平行世界。创建的Checkpoint的标识会存储返回的RunnableConfig配置中。1. 以节点名义更新状态上述这些方法总是以一个节点的名义模拟一个具体的任务来更新对应的状态值。这个模拟的任务不会真正被执行必须关联一个有效的节点因为需要借助于对应节点的writers列表完成写入操作。针对单一任务的更新请求通过一个StateUpdate对象来表示它的as_node和task_id字段表示状态更新的模拟任务的节点名称和任务标识。classPregel(PregelProtocol[StateT,ContextT,InputT,OutputT],Generic[StateT,ContextT,InputT,OutputT]):defbulk_update_state(self,config:RunnableConfig,supersteps:Sequence[Sequence[StateUpdate]],)-RunnableConfigasyncdefabulk_update_state(self,config:RunnableConfig,supersteps:Sequence[Sequence[StateUpdate]],)-RunnableConfigdefupdate_state(self,config:RunnableConfig,values:dict[str,Any]|Any|None,as_node:str|NoneNone,task_id:str|NoneNone,)-RunnableConfigasyncdefaupdate_state(self,config:RunnableConfig,values:dict[str,Any]|Any,as_node:str|NoneNone,task_id:str|NoneNone,)-RunnableConfigclassStateUpdate(NamedTuple):values:dict[str,Any]|Noneas_node:str|NoneNonetask_id:str|NoneNone2. 单一节点更新update_state/aupdate_state方法用于单一节点的状态更新。如果没有利用as_node参数显式指定节点这两个方法会默认使用最后一次实施更新的节点。那么如何确定最后一次状态更新来源于哪个节点呢还记得Checkpoint的channel_versions和versions_seen字段吗前者返回所有通道的最新版本后者返回每个节点可见的通道和版本。能够看到最高版本的那个通道的所有节点都会被认为是实施了最新的状态更新。如果通过这种策略解析出来的节点不止一个就会抛出异常并提醒我们提供确定的节点如下的演示程序体现了这一点。fromlanggraph.channelsimportLastValuefromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.checkpoint.memoryimportInMemorySaverfromfunctoolsimportpartialdefhandle(node:str,args:dict)-str:returnnode foo(NodeBuilder().subscribe_to(start,readFalse).do(partial(handle,foo)).write_to(bar))bar1(NodeBuilder().subscribe_to(bar,readFalse).do(partial(handle,bar1)).write_to(bar1))bar2(NodeBuilder().subscribe_to(bar,readFalse).do(partial(handle,bar2)).write_to(bar2))appPregel(nodes{foo:foo,bar1:bar1,bar2:bar2},channels{start:LastValue(None),bar:LastValue(str),bar1:LastValue(str),bar2:LastValue(str),},input_channels[start],output_channels[bar1,bar2],checkpointerInMemorySaver(),)config{configurable:{thread_id:tx123}}resultapp.invoke(input{start:None},configconfig)assertresult{bar1:bar1,bar2:bar2}new_configapp.update_state(configconfig,values{bar1:bar1[new]},)stateapp.get_state(confignew_config)print(state.values)上面这个Pregel体现的执行流程很简单启动的时候通过写入通道start驱动节点foo的执行后者结束之后写入通道bar驱动节点bar1和bar2的执行这两个节点最终会将自身的节点名称写入对应的通道bar1和bar2。在完成常规调用之后我们调用update_state方法试图将最终通道bar1的值改写成bar1[new]但是最终会抛出InvalidUpdateError异常并提示Ambiguous update, specify as_node。由于没有在作为输入参数的RunnableConfig配置中指定Checkpoint的ID所以针对状态的修改是基于最后创建的Checkpoint进行的。按照我们上面说的逻辑引擎会根据每个节点可见通道的最高版本得到最近一次更新状态的节点。很明显节点bar1和bar2订阅的通道bar具有更高的版本。由于不满足节点的唯一性所以InvalidUpdateError异常被抛出。由于通道bar1是由节点bar1写入的所以我们将as_node参数设置为bar1实际上你设置任意的节点名称都可以只要对应节点存在即可。new_configapp.update_state(configconfig,values{bar1:bar1[new]},as_nodebar1)stateapp.get_state(confignew_config)print(state.values)你以为这就结束了吗虽然这次状态更新成功了但是最新的状态却非我所愿。从如下的输出结果可以看出bar1的值是{bar1: bar1[new]}也就是说它是将作为参数values的字典整个作为了写入的值。{start:None,bar:foo,bar1:{bar1:bar1[new]},bar2:bar2}要解释这个问题就必须真正了解状态究竟是如何被更新的。我们之所以需要确定以哪个节点的名义更新状态并不仅仅是为了补充必要的审计信息而是因为整个更新操作依赖于对应PregelNode的writers列表。我们再回顾一下PregelNode如下所示的writers字段它返回的通道写入器体现为一组Runnable对象。classPregelNode:writers:list[Runnable]cached_propertydefflat_writers(self)-list[Runnable]classChannelWrite(RunnableCallable):writes:list[ChannelWriteEntry|ChannelWriteTupleEntry|Send]def__init__(self,writes:Sequence[ChannelWriteEntry|ChannelWriteTupleEntry|Send],*,tags:Sequence[str]|NoneNone,)classChannelWriteEntry(NamedTuple):channel:strvalue:AnyPASSTHROUGH skip_none:boolFalsemapper:Callable|NoneNoneclassChannelWriteTupleEntry(NamedTuple):mapper:Callable[[Any],Sequence[tuple[str,Any]]|None]value:AnyPASSTHROUGH static:Sequence[tuple[str,Any,str|None]]|NoneNonePASSTHROUGHobject()每个通道写入器对应一个ChannelWrite对象后者针对通道的写入意图会被添加到writes字段对用的列表中这是由一组ChannelWriteEntry、ChannelWriteTupleEntry或者Send对象的列表。当update_state/aupdate_state方法将节点确定下来后这个列表被提取出来对于ChannelWriteEntry和ChannelWriteTupleEntry其values字段被替换成传入update_state/aupdate_state方法的values参数仅此而已。这就是我们提供的字典作为整体被写入通道bar1的原因所以调用update_state/aupdate_state方法的时候按照如下的方式提供具体的值就好。new_configapp.update_state(configconfig,valuesbar1[new],as_nodebar1)我们在前面说过update_state/aupdate_state方法是通过在最新或者指定Checkpoint基础上创建一个新的Checkpoint进而达到更新状态的目的我们可以输出完整的历史来证明这一点。forstateinapp.get_state_history(config):metadatastate.metadata stepmetadata[step]sourcemetadata[source]print(fstep:{step}\nsource:{source}\nvalues:{state.values}\n)在完成状态更新后我们使用上面的代码提取组成历史的每个快照并将对应的Superstep编号、source和values输出来。在如下所示的输出中Superstep 2对应的快照就是调用update_state方法产生的这个快照还具有不同的sourceupdate揭示它的与众不同。step:2 source:update values: {start: None, bar: foo, bar1: bar1[new], bar2: bar2} step:1 source:loop values: {start: None, bar: foo, bar1: bar1, bar2: bar2} step:0 source:loop values: {start: None, bar: foo} step:-1 source:input values: {start: None}3. 更新失效如果我们对上述的更新原理不了解在遇到一些状态更新失效的场景时可能永远找不到问题的症结。比如在如下这个简单的例子中Pregel唯一的节点会将值foo写入通道调用之后针对结果的断言也证实了写入时成功的。fromlanggraph.channelsimportLastValuefromlangchain_core.runnablesimportRunnableConfigfromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.checkpoint.memoryimportInMemorySaver node(NodeBuilder().subscribe_only(foo).do(lambdaargs:args).write_to(outputlambda_:foo))appPregel(nodes{node:node},channels{foo:LastValue(str),output:LastValue(str),},input_channels[foo],output_channels[output],checkpointerInMemorySaver(),)config:RunnableConfig{configurable:{thread_id:tx123}}resultapp.invoke(input{foo:foo},configconfig)assertresult[output]foonew_configapp.update_state(configconfig,valuesbar,as_nodenode)stateapp.get_state(new_config)assertstate.values[output]foo# state remains unchanged我们本希望调用update_state方法将输出改写为bar。但是从调用get_state的结果来看状态并没有更新成功。那么是因为新的Checkpoint没有创建吗为此我们按照如下的方式输出整个历史。forstateinapp.get_state_history(config):metadatastate.metadata stepmetadata[step]sourcemetadata[source]print(fstep{step}\nsource:{source}\nvalues:{state.values})print()从如下的输出结果可以看出update_state方法调用对应的Checkpoint已经成功创建但是它的状态values就是没有改变。step 1 source: update values: {start: None, output: foo} step 0 source: loop values: {start: None, output: foo} step -1 source: input values: {start: None}这个的问题出在节点的构建上面由于我们调用调用NodeBuilder的write_to方法时采用了关键字参数来确定目标通道并以Lambda表达式的方式提供写入的值并且Lambda表达式并没有使用原始输入而是直接硬编码成foo。这行代码将生成一个ChannelWriteEntry对象并为它指定正确的通道名称output它的value不会被设置为处理函数的返回值但是Lambda表达式转换成的Callable对象将作为mapper字段。由于update_state方法仅仅通过对节点的写入器稍加改造来完成通道写入。对于这个背后创建的ChannelWriteEntry对象来说改变的只有其value字段mapper字段将保持不变。由于mapper无脑返回foo所以状态永远也不可能被改变。3. 对PendingWrite的影响前面介绍的更新都是在一个不存在PendingWrite的状态上完成的如果Superstep尚未完结并且同时具有完成和中断的任务update_state方法调用后整个状态又是什么样子呢经过我的测试不论我们以完成任务对应的节点的名义还是以中断任务对应的节点的名义针对update_state方法的调用都会以创建新的Checkpoint的方式强行闭合当前Superstep。以如下这段程序为例Pregel的两个并行执行的初始节点foo和bar前者成功执行后者会遇到中断。我们调用update_state方法以节点bar的名义试图更新通道bar的状态。在调用前后我们输出整个历史。fromlanggraph.channelsimportLastValuefromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.checkpoint.memoryimportInMemorySaverfromlanggraph.typesimportinterruptdefhandle(args:dict)-str:resumeinterrupt(Resuming execution)returnresume foo(NodeBuilder().subscribe_to(start,readFalse).do(lambdaargs:args).write_to(foo))bar(NodeBuilder().subscribe_to(start,readFalse).do(handle).write_to(bar))appPregel(nodes{foo:foo,bar:bar},channels{start:LastValue(None),foo:LastValue(str),bar:LastValue(str),},input_channels[start],output_channels[foo,bar],checkpointerInMemorySaver(),)config{configurable:{thread_id:tx123}}defshow_history(config):for_,checkpoint,metadata,_,pending_writesinapp.checkpointer.list(config):stepmetadata[step]sourcemetadata[source]print(fstep:{step}\nsource:{source}\nvalues:{checkpoint[channel_values]}\npending_writes:{pending_writes}\n)app.invoke(input{start:None},configconfig)print(After invoke:)show_history(config)new_configapp.update_state(configconfig,valuesupdated value,as_nodebar)print(\nAfter update_state:)show_history(config)从如下的输出结果可以看出正常调用之后确实遇到了中断。update_state方法调用之后通道bar的状态确实被更新。但是新的Checkpoint被创建后之前的PendingWrite将不复存在。After invoke: step:-1 source:input values: {start: None} pending_writes: [(a2357188-2d04-182f-8672-0328182f68a0, __interrupt__, [Interrupt(valueResuming execution, idfcb47fce081f41d1e1141d004a9b66f9)]), (10b0e7ea-3453-e5ad-18a7-eb47f3caf108, foo, {})] After update_state: step:0 source:update values: {start: None, bar: updated value} pending_writes: [] step:-1 source:input values: {start: None} pending_writes: [(a2357188-2d04-182f-8672-0328182f68a0, __interrupt__, [Interrupt(valueResuming execution, idfcb47fce081f41d1e1141d004a9b66f9)]), (10b0e7ea-3453-e5ad-18a7-eb47f3caf108, foo, {}), (a2357188-2d04-182f-8672-0328182f68a0, bar, updated value)]本例中我们是以中断节点的名义对状态实施修改的所以PendingWrite被新的状态抹除还说得过去的。但是如果我们按照如下的方式以成功执行的节点foo的名义修改通道foo的状态呢new_configapp.update_state(configconfig,valuesupdated value,as_nodefoo)从如下所示的输出结果可以看出虽然我们对中断节点bar没有实施任何操作update_state方法依然会将其PendingWrite抹除。我们从这个例子大体可以看出update_state背后的逻辑针对最终状态的更新总是在最新的Checkpoint描述的状态下进行并且创建一个新的Checkpoint作为最终的状态。After invoke: step:-1 source:input values: {start: None} pending_writes: [(87a80b8a-6c7b-a026-68b5-dd39d405a6a4, foo, {}), (8a53715a-dc21-56ce-d0a6-d13ff6604d15, __interrupt__, [Interrupt(valueResuming execution, id8672a7e5d7233e1f8c11fc9ff4b00c56)])] After update_state: step:0 source:update values: {start: None, foo: updated value} pending_writes: [] step:-1 source:input values: {start: None} pending_writes: [(87a80b8a-6c7b-a026-68b5-dd39d405a6a4, foo, {}), (8a53715a-dc21-56ce-d0a6-d13ff6604d15, __interrupt__, [Interrupt(valueResuming execution, id8672a7e5d7233e1f8c11fc9ff4b00c56)])]
返回列表