尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
[拆解LangChain执行引擎-04]ManagedValue:一种特殊的只读虚拟通道
我们一直在强调Pregel对象的状态是通过通道维护和传递的其实承载传递状态功能的除了通道还有ManagedValue。我们可以将ManagedValue视为虚拟通道节点不仅采用与读取通道完全一样的方式读取ManagedValue而且注册的ManagedValue也直接存放在Pregel的channels字段中。1. ManagedValue的定义如果我们仔细查看Pregel类的定义可以看出其channels字段返回一个字典字典的值的类型BaseChannel和ManagedValueSpec两种类型的联合前者是通道的基类后者就是ManagedValue类的别名。classPregel(PregelProtocol[StateT,ContextT,InputT,OutputT],Generic[StateT,ContextT,InputT,OutputT]):channels:dict[str,BaseChannel|ManagedValueSpec]ManagedValueSpectype[ManagedValue]如果说通道存储的是的业务状态那么ManagedValue传递的就是Pregel这个执行引擎的运行时状态。一般来说ManagedValue自身不负责存储状态因为它提供的值可以实时计算出来自然也不参与持久化。从如下所示的代码片段可以看出ManagedValue仅仅定义了一个唯一的静态抽象方法get返回对应的值由于作为输入的PregelScratchpad对象提供的信息有限所以ManagedValue能够发挥的空间其实不大在大部分情况下根本用不到它。classManagedValue(ABC,Generic[V]):staticmethodabstractmethoddefget(scratchpad:PregelScratchpad)-V:...2. PregelScratchpadManagedValue提供的值是通过静态get方法根据PregelScratchpad对象计算所得。当确定后续待执行的节点后引擎会为每个节点创建一个任务每个任务都会附加一个PregelScratchpad对象。PregelScratchpad的step和stop字段就返回当前Superstep的编号和允许的最大迭代的次数相当Superstep编号的最大值其它字段与持久化有关。dataclasses.dataclass(**_DC_KWARGS)classPregelScratchpad:step:intstop:intcall_counter:Callable[[],int]interrupt_counter:Callable[[],int]get_null_resume:Callable[[bool],Any]resume:list[Any]subgraph_counter:Callable[[],int]PregelScratchpad的call_counter、interrupt_counter和subgraph_counter字段以闭包的形式返回一个计数器。call_counter计数器用于为当前Superstep内产生的所有任务分配唯一的内部序列号。2.1 Resume Value和中断计数器interrupt_counter、get_null_resume和resume字段与Pregel基于中断Interrupt/恢复Resume的执行方式有关。假设Pregel的对应一个需要人工介入的多级审批流程在每次需要人工介入收集审批者决定的时候流程进入一个中断当前的状态被持久化。当用户提供审批结果后流程以恢复的形式执行此时中断时持久化的快照被提取出来恢复现场审批结果以Resume Value的形式提供给Pregel。如果涉及多轮中断为了匹配每个中断点与对应的Resume Value之间的映射关系提供的Resume Value会按照顺序被持久化并在恢复执行的时候与当前当前提供的Resume Value进行合并后一并填充到PregelScratchpad的resume列表中。恢复执行无法真正做到在中断点出开始执行它只能从头执行节点的处理函数所以定义幂等节点函数应该成为Agent编程的金科玉律。由于PregelScratchpad的resume字段会按照中断的顺序存放Resume Value所以在恢复执行的时候每遇到一个中断引擎可以利用interrupt_counter字段返回的计数器作为索引从resume列表中将匹配的Resume Value提取出来。如果提取的Resume Value为None或者计数器返回的索引越界get_null_resume字段提供的回调函数就会执行。这个回调函数具有一个bool类型的参数is_called调用时该参数被设置为True表示该中断确实被触发了但没有对应的数据。这会消耗掉这个中断位确保流程不至于永远得不到恢复。2.2 子图调用计数器如果说interrupt_counter计数器旨在解决每次中断与提供Resume Value的匹配问题那么subgraph_counter计数器解决的每次子图调用与对应Pregel实例的匹配问题。如果站在图的视角每个Pregel对象就是由多个节点组成的图而Pregel可以作为一个节点出现在另一个Pregel构建的图中两个Pregel之间就称为了父子关系子Pregel构建的图就是子图针对它的调用就是子图调用。虽然在同一个图中但是每个Pregel会独自完成自身的持久化。在恢复执行场景中引擎会率先加载作为根Pregel对应的Checkpoint来恢复现场。当遇到以子图形式调用另一个Pregel时引擎会加载对应的Checkpoint来恢复目标Pregel在那个时刻的状态。现在问题来了在目标Pregel众多持久化的Checkpoint中怎么知道该加载哪一个呢这个问题本质上描述的是当某个Pregel作为子图被调用时持久化生成的Checkpoint如何与当前执行上下文进行匹配这个问题可用利用Checkpoint的命名空间来解决的那么命名空间由哪些要素组成呢我们知道节点是以任务的形式被执行的每个任务具有唯一的ID并且在恢复时保持不变如果命名空间由执行链路上每个任务的节点名称任务ID组成那么子图对应Pregel的Checkpoint就能利用此命名空间关联起来。但是问题还是没有完全解决如果同一个任务涉及针对同一Pregel的多次调用如命名空间只包含基于任务的执行路径此时两个子图会共享相同的命名空间具体对应哪个Checkpoint依然无法解决所以Checkpoint的命名空间还应该包含调用序号。Checkpoint的命名空间的规则可以通过如下这个演示实例来证实。如代码片段所示我们创建了一个由单一节点组成的Pregel对象sub_graph命名为baz的节点在执行的时候会从当前的RunnableConfig配置中提取并输出当前的Checkpoint命名空间。fromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.channelsimportLastValuefromlanggraph.checkpoint.memoryimportInMemorySaverfromlanggraph.pregel._writeimportChannelWrite,ChannelWriteTupleEntryfromlanggraph.typesimportRunnableConfigfromtypingimportAnydefhandle(args:dict[str,Any],config:RunnableConfig)-None:print(config[configurable][checkpoint_ns])sub_node(NodeBuilder().subscribe_to(start).do(handle))sub_graphPregel(nodes{baz:sub_node},channels{start:LastValue(None),},input_channels[start],output_channels[])defhandle1(args:dict[str,Any])-None:sub_graph.invoke(input{start:None})defhandle2(args:dict[str,Any])-str:sub_graph.invoke(input{start:None})sub_graph.invoke(input{start:None})foo(NodeBuilder().subscribe_to(foo).do(handle1).write_to(barNone))bar(NodeBuilder().subscribe_to(bar).do(handle2))graphPregel(nodes{foo:foo,bar:bar},channels{foo:LastValue(None),bar:LastValue(str),},input_channels[foo],output_channels[],checkpointerInMemorySaver())config{configurable:{thread_id:123}}graph.invoke(input{foo:None},configconfig)在另一个Pregel中我们为它设置了两个先后执行的节点foo和bar前者调用sub_graph一次后者调用两次。针对三次调用sub_graph为自身持久化设置的Checkpoint命名空间会以如下的形式输出可以看出命名空间同时体现了调用链路和次序。foo:36817c76-c3f7-643f-7924-0d29b39f469a|baz:311cc911-96a0-56b6-225b-28e4cece7cd9 bar:97be6a71-1b71-7364-e691-a122cfef1a92|baz:789287de-869f-42b8-dd03-7518820daaa6 bar:97be6a71-1b71-7364-e691-a122cfef1a92|1|baz:dd1ddd1b-fc62-b46a-c2ec-6a1d8344b793由于Pregel支持基于中断/恢复的执行方式我们应该重新思考Pregel实例这个概念。即使程序中先后指定的两个地方引用的是同一个变量如果中间出现中断它们引用的就不会是同一个Pregel实例。在不断的中断/恢复执行流程中所谓Pregel实例有时候表示成对应的Checkpoint可能更准确。对于同一个节点任务来说如果涉及针对同一个子Pregel的多次调用从第二次调用开始对方持久化生成的Checkpoint会将调用次序包含在命名空间中。在恢复执行的时候自然也需要根据当前的执行上下文提供包含此序号的命名空间采用加载对应的Checkpoint并最终恢复对应的Pregel对象PregelScratchpad的subgraph_counter字段返回的计数器就是为了提供这个序号。3. 两个原生的ManagedValue由于ManagedValue所能提供的值是根据PregelScratchpad计算生成而后者可用的唯有表示当前和最大Superstep编号的step和stop字段所以我们采用ManagedValue的应用场景其实很优先。我从只找到如下两个原生的ManagedValue类型它们都定义在langgraph.managed.is_last_step这个包中。其中一个IsLastStepManager用于判断是否为最后一个Superstep而RemainingStepsManager则用来确定余下的Superstep数。具体的实现非常简单仅仅是针对PregelScratchpad的step和stop字段的简单运算而已。classIsLastStepManager(ManagedValue[bool]):staticmethoddefget(scratchpad:PregelScratchpad)-bool:returnscratchpad.stepscratchpad.stop-1classRemainingStepsManager(ManagedValue[int]):staticmethoddefget(scratchpad:PregelScratchpad)-int:returnscratchpad.stop-scratchpad.step由于ManagedValue属于一个计算属性所以它只能作为节点的输入节点针对ManagedValue和常规通道的读取方式完全一致。在创建Pregel对象时所用到的ManagedValue需要在channels字段中显式声明但是不能将其添加到输入和输出通道列表中。如下的实例演示了RemainingStepsManager的使用方式。如代码所示我们创建的Pregel由两个先后执行的节点构成foo和bar它们会将命名为remaining_steps的ManagedValue作为输入并将其分别输出到remaining_steps_after_foo和remaining_steps_after_bar这两个通道中分别表示在这两个节点完成执行后所剩的Superstep数。fromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.managed.is_last_stepimportRemainingStepsManagerfromlanggraph.channelsimportLastValue foo(NodeBuilder().subscribe_to(foo).read_from(remaining_steps).do(lambdaargs:args[remaining_steps]).write_to(remaining_steps_after_foolambdaargs:args,barNone))bar(NodeBuilder().subscribe_to(bar).read_from(remaining_steps).do(lambdaargs:args[remaining_steps]).write_to(remaining_steps_after_bar))appPregel(nodes{foo:foo,bar:bar},channels{foo:LastValue(None),bar:LastValue(None),remaining_steps_after_foo:LastValue(int),remaining_steps_after_bar:LastValue(int),remaining_steps:RemainingStepsManager,},input_channels[foo],output_channels[remaining_steps_after_foo,remaining_steps_after_bar])config{recursion_limit:10}resultapp.invoke({foo:None},configconfig)assertresult[remaining_steps_after_foo]10assertresult[remaining_steps_after_bar]9在根据两个节点创建Pregel对象时我们将针对命名为remaining_steps的ManagedValue的声明添加到channels字段中对应的类型被设置为RemainingStepsManager。由于在调用Pregel对象时利用RunnableConfig配置将Superstep迭代限制为10所以先后执行的两个节点后剩余步数分别为10和9。
RELATED

相关推荐

[拆解LangChain执行引擎-03]__pregel_tasks通道:成就“PUSH任务”的功臣

[拆解LangChain执行引擎-03]__pregel_tasks通道:成就“PUSH任务”的功臣

除了我们显式声明的用于存储业务数据或驱动信号的通道之外,Pregel自身也会维护一些系统通道,其中最重要的莫过于一个名为__pregel_tasks的通道。通过前面针对BSP的介绍,我们知道当Superstep进入同步屏障并应用所有更新后,引擎会根…

📅 2026/10/3 2:06:33
工业日志结构化与PDF表格提取:Profinet/Modbus数据解析实战

工业日志结构化与PDF表格提取:Profinet/Modbus数据解析实战

从抓包到表格:工业日志结构化与PDF提取的完整实操记录搞工业数据处理的人,大概都经历过那种“数据在眼前,就是拿不到”的崩溃感。明明PLC就在机房里闪着灯,传感器数据一条条往上传,但当你打开抓包文件或者翻看设备导出…

📅 2026/10/3 2:06:33
[拆解LangChain执行引擎-08]基于Checkpoint的持久化

[拆解LangChain执行引擎-08]基于Checkpoint的持久化

LangGraph基于Checkpoint的持久化,核心是在每个Superstep后保存Pregel的状态,并以thread_id组织为可恢复轨迹。其作用包括支持容错与断点续跑,崩溃或中断后从最近检查点恢复;支撑人机交互,暂停等待审批后继续&#xff…

📅 2026/10/3 2:06:33
MORE NEWS

更多资讯

📰

用 thiserror 派生宏消除自定义错误样板代码:100-exercises-to-learn-rust 的 TicketNewError 实战

示例工程教程 【免费下载链接】100-exercises-to-learn-rust A self-paced course to learn Rust, one exercise at a time. 项目地址: https://gitcode.com/GitHub_Trending/10/100-exercises-to-learn-rust 点击查看 免费下载 本篇指南聚焦 Rust 生态中最常用的错…

📰

IDM-VTON 人体解析工具链:Detectron2 tools 目录训练、评测与可视化脚本全解析

计算机视觉深度学习媒体生成 【免费下载链接】IDM-VTON [ECCV2024] IDM-VTON : Improving Diffusion Models for Authentic Virtual Try-on in the Wild 项目地址: https://gitcode.com/GitHub_Trending/id/IDM-VTON 点击查看 免费下载 导读:本文围绕 I…

📰

基于多视觉语言模型交叉描述的智能眼镜图像理解与质量评估实战指南(OpenGlass 项目)

人工智能AI 应用智能硬件本地部署可穿戴AI Agent 【免费下载链接】OpenGlass Turn any glasses into AI-powered smart glasses 项目地址: https://gitcode.com/GitHub_Trending/op/OpenGlass 点击查看 免费下载 OpenGlass 是一个让任何普通眼镜变身 AI 智能眼镜的…

📰

Toonflow 更新说明全解读:从 21 种语言界面到画布复制、FFmpeg 工具与桌面更新机制

人工智能AI 应用AI AgentRAGAI 写作后端桌面应用 【免费下载链接】Toonflow-app Toonflow 是一款 AI 短剧漫剧工具,能够利用 AI 技术将小说自动转化为剧本,并结合 AI 生成的图片和视频,实现高效的短剧创作。借助 Toonflow,可以轻松…

📰

AI-For-Beginners Game Jam 作业实战指南:以「过去—现在—未来」框架剖析游戏中的 AI 进化

教程人工智能机器学习深度学习 【免费下载链接】AI-For-Beginners 12 Weeks, 24 Lessons, AI for All! 项目地址: https://gitcode.com/GitHub_Trending/ai/AI-For-Beginners 点击查看 免费下载 本指南基于 AI-For-Beginners 第 1 课(Introduction to A…

📰

telegram - api-reference

Telegram Bot API - 完整参考 目录 认证发送方法编辑方法聊天方法成员方法更新与 Webhook机器人配置主要类型解析模式错误代码 认证 基础 URL&#xff1a; https://api.telegram.org/bot<TOKEN>/<METHOD> 文件 URL&#xff1a; https://api.telegram.org/file/b…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬