[拆解LangChain执行引擎-10]非常规PendingWrite的持久化

发布时间:2026/10/4 12:25:56
[拆解LangChain执行引擎-10]非常规PendingWrite的持久化 持久化状态的提取重点介绍介绍用了PendingWrite这个三元组它用于描述在一个尚未完成的Superstep中的某个节点针对通道的写入。这个三元组的第二部分表示写入的通道。对于一些特殊的场景比如出错、无写入、中断和恢复它们的值不再是一个普通的通道名称而是使用如下的值__error__执行节点任务出现异常__no_writes__节点任务成功执行但是没有执行针对通道的输出__interrupt__任务中断__resume__表示恢复执行提供的数据。1. 抛出异常接下来我们两个例子来产生上述这几种特殊的PendingWrite。我们先来模拟出错的场景如下面的代码片段所示我们执行的Pregel对象具有一个唯一的节点它的处理函数直接抛出一个异常。fromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.channelsimportLastValuefromlanggraph.checkpoint.memoryimportInMemorySaverfromtypingimportAnydefhandle(args:dict[str,Any])-None:raiseException(manllually raised exception)nodeNodeBuilder().subscribe_to(start).do(handle)appPregel(nodes{body:node},channels{start:LastValue(str)},checkpointerInMemorySaver(),input_channels[start],output_channels[],)config{configurable:{thread_id:123}}try:resultapp.invoke({start:begin},configconfig)exceptExceptionasex:print(fCaught exception:{ex})(_,_,_,_,pending_writes)app.checkpointer.get_tuple(config)print(pending_writes)我们在try/except块中完成针对Pregel的调用并捕捉和输出得到的异常信息。接下来我们调用Checkpointer一个InMemorySaver对象的get_tuple方法得到对应的CheckpointTuple元组然后将pending_writes部分输出。从如下所示的输出结果可以看出这个PendingWrite三元组的通道名称被设置为__error__整个Exception对象成为了写入的内容。Caught exception:manllually raised exception [(f9ff1e88-4d82-f417-ad11-8fd870bfe647, __error__, Exception(manllually raised exception))]2. 无通道输入、中断与恢复并不是所有的节点都有向通道写入执行结果的需求只要处理函数成功执行即使没有通道输出的行为该任务的状态也会被视为成功持久化时会采用不同的形式来记录这种不需要写入的PendingWrite。如下的这个程序不仅仅演示了这种无输出写入的场景还同时模拟了中断和恢复。fromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.channelsimportLastValuefromlanggraph.checkpoint.memoryimportInMemorySaverfromtypingimportAnyfromlanggraph.typesimportCommand,interruptdeffoo(args:dict[str,Any])-list[str]:resume1interrupt(1st interrupt)assertresume11st resumeresume2interrupt(2nd interrupt)assertresume22nd resumeresume3interrupt(3rd interrupt)assertresume33rd resumereturn[resume1,resume2,resume3]defbar(args:dict[str,Any])-None:passappPregel(nodes{foo:NodeBuilder().subscribe_only(start).do(foo).write_to(output),bar:NodeBuilder().subscribe_only(start).do(bar),},channels{start:LastValue(str),output:LastValue(list[str]),},input_channels[start],output_channels[output],checkpointerInMemorySaver(),)config{configurable:{thread_id:123}}resultapp.invoke(input{start:begin},configconfig,stream_modetasks)(_,_,_,_,pending_writes)app.checkpointer.get_tuple(config)print(fAfter invoke:\n{pending_writes})app.invoke(inputCommand(resume1st resume),configconfig)(_,_,_,_,pending_writes)app.checkpointer.get_tuple(config)print(f\nAfter resume 1:\n{pending_writes})app.invoke(inputCommand(resume2nd resume),configconfig)(_,_,_,_,pending_writes)app.checkpointer.get_tuple(config)print(f\nAfter resume 2:\n{pending_writes})resultapp.invoke(inputCommand(resume3rd resume),configconfig)assertresult{output:[1st resume,2nd resume,3rd resume]}(_,_,_,_,pending_writes)app.checkpointer.get_tuple(config)print(f\nAfter resume 3:\n{pending_writes})如上面的代码片段所示我们为Pregel提供了两个并行执行的节点foo和bar其中bar对应的函数并未执行任何有效操作也没有任何的输出。我们为节点foo对应的处理函数制造了三次人为中断所以整个流程需要至少四次调用才能结束。我们在创建的RunnableConfig对象中提供了统一的thread_id并将它作为后续方法调用的参数。针对Pregel的三次调用第一次是为常规调用后面两次分别是针对两次中断的恢复调用。我们在每次调用后得到并输出持久化记录下来的PendingWrite。从如下的输出结果可以看出第一次常规调用后节点foo停在第一个中断处节点bar成功执行但没有输出这种状态以PendingWrite记录下来通道名称分别是__interrupt__和__no_writes__。前者的写入内容是一个Interrupt对象它具有我们指定的值1st interrupt。我们也看到了Interrupt对象具有一个唯一标识在恢复调用时我们可以利用此标识为其指定针对性的恢复数据Command(resume{id:resume value)。After invoke: [(8d407c25-02f6-9101-d1b8-5a99c247edde, __interrupt__, [Interrupt(value1st interrupt, id5603cdf275d8b8ba0633d272fa176fd3)]), (22507855-e257-1b5b-eb1a-3c3fb0a071e9, __no_writes__, None)] After resume 1: [(8d407c25-02f6-9101-d1b8-5a99c247edde, __interrupt__, [Interrupt(value2nd interrupt, id5603cdf275d8b8ba0633d272fa176fd3)]), (22507855-e257-1b5b-eb1a-3c3fb0a071e9, __no_writes__, None), (00000000-0000-0000-0000-000000000000, __resume__, 1st resume), (8d407c25-02f6-9101-d1b8-5a99c247edde, __resume__, [1st resume])] After resume 2: [(8d407c25-02f6-9101-d1b8-5a99c247edde, __interrupt__, [Interrupt(value3rd interrupt, id5603cdf275d8b8ba0633d272fa176fd3)]), (22507855-e257-1b5b-eb1a-3c3fb0a071e9, __no_writes__, None), (00000000-0000-0000-0000-000000000000, __resume__, 2nd resume), (8d407c25-02f6-9101-d1b8-5a99c247edde, __resume__, [1st resume, 2nd resume])] After resume 3: []针对第一个中断的恢复调用后节点foo停在第二个中断处此时会创建两个新的PendingWrite持久化我们提供的Resume Value1st resume此时采用的通道名称就是__resume__但为什么是两个呢这实际上反映了Pregel 处理外部指令注入与节点内部消费的同步机制。第一个被称为全局Resume ValueGlobal Resume Value, 它代表从外部通过Command(resume...)注入到图中的原始指令。由于它不是由图内节点产生的因此task_id为空它是唤醒整个暂停状态的总开关。第二个节点foo对全局Resume Value的消费记录所以具有一个明确的task_id。当节点foo被唤醒并执行到interrupt行时它会从全局Resume Value读取数据。为了保证幂等性和可回溯性系统会将拿走了哪个Resume Value记录在它的任务路径下。针对Resume的冗余设置是为了解决重入与回溯问题。全局记录证明了用户确实提供了这个值。节点记录证明了这个值确实被这个特定的interrupt函数调用消费了。一个节点内部可以连续调用多次interrupt函数系统需要按顺序记录该节点消费过的所有Resume Value以便在时间旅行或恢复调用时能够精确对齐。当我们调用interrupt函数实施人为中断时底层实际上会抛出一个GraphInterrupt异常Pregel通过捕获这个异常进而生成针对性的PendingWrite所以针对同一个任务只可能有一个此类中断类型的PendingWrite。由于恢复执行总是会从头执行节点函数所以基于中断的PendingWrite并不会对恢复执行造成任何影响。当我们完成第二次恢复调用后持久化的中断PendingWrite反映的是针对第二次interrupt函数的调用对应Interrupt对象的值为2nd interrupt。Resume Value必须按照顺序提供因为每遇到一个interrupt函数的调用都会利用前面介绍过的计算器提供的索引从Resume Value列表中读取对应的值作为该调用的返回值所以持久化的第二个基于恢复的PendingWrite对应的值变成了包含两个Resume Value的列表[1st resume, 2nd resume]。在针对第三个中断的恢复执行结束后foo节点完成了它的执行任务而bar节点对应的任务本身就是成功状态所以整个Superstep顺利结束自然也就不存在PendingWrite了。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询