恒美微站 Logo 恒美微站
  • 首页
  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心
  • 联系我们

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

  • 首页
  • 资讯中心
  • /
  • [拆解LangChain执行引擎-03]__pregel_tasks通道:成就“PUSH任务”的功臣

相关资讯

工业日志结构化与PDF表格提取:Profinet/Modbus数据解析实战 2026/10/3 2:06:29
[拆解LangChain执行引擎-08]基于Checkpoint的持久化 2026/10/3 2:06:29
fantastic-admin 标签栏(Tabbar)配置详解:icon 与 hotkeys 实战指南 2026/10/3 2:06:29

最新资讯

K8s集群Kubectl命令高阶运维实操
从Postman到Apifox:API一站式协作与自动化测试实战指南
Apifox从入门到实战:接口调试、Mock与自动化测试全攻略
分布式锁选型指南:Redis、ZooKeeper与数据库方案全解析
分布式锁选型指南:Redis、ZooKeeper与数据库锁方案实战对比
提示工程+LoRA微调:让大模型生成可直接进CI的Java单元测试

今日推荐

SAP生产预留实战指南:MB21/MB23/MB25协同与MRP集成
编译原理实验:递归下降分析器消除左递归与避坑指南
Python协议级爬取Shopee商品数据实战

本周热门

从像素到笔画:srt-whiteboard-animation骨架笔迹追踪实现(Zhang-Suen细化+8邻接追踪)
网站建设的英语怎么说?别只背单词,看完这套安全完整流程才敢上线
新手入门看这篇:建设网站加盟避坑指南与SEO实操

本月精选

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证
2026 大模型集体涨价:用 Python 做企业 Token 成本测算与选型避坑(附配置)

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

发布时间:2026/10/3 2:06:29
[拆解LangChain执行引擎-03]__pregel_tasks通道:成就“PUSH任务”的功臣 除了我们显式声明的用于存储业务数据或驱动信号的通道之外Pregel自身也会维护一些系统通道其中最重要的莫过于一个名为__pregel_tasks的通道。通过前面针对BSP的介绍我们知道当Superstep进入同步屏障并应用所有更新后引擎会根据节点针对通道的订阅和通道自身的是否发生改变生成下一步待执行的任务其实待执行的任务的生成方式不限于此。1. 两种任务创建方式我们将根据节点针对通道的订阅来驱动任务执行的模式称为Pull模式与之相对的则是借助于__pregel_tasks这个通道实现的Push模式。这是一个关闭累积模式的Topic类型的通道它存储的Topic体现为具有如下定义的Send对象。当某个节点执行之后可以像这个通道中写入一个Send来驱动某个节点在下一Superstep中执行。除了利用Send对象的node字段指定待执行的节点名称外还可以利用arg字段提供输入参数。classSend:node:strarg:Anydef__init__(self,/,node:str,arg:Any)-None由于关闭了累积模式在Topic类型通道中写入的内容只会在下一个Superstep中生效并且具有阅后即焚的特性。对于执行引擎来说这个名为__pregel_tasks的通道存储的就是下一Superstep以Push模式驱动执行的任务列表两者完美契合。2. 确认__pregel_tasks通道的存在__pregel_tasks通道的存在可以通过如下的演示实例来验证。如代码片段所示在采用常规方式将Pregel对象创建出来后我们根据通道名称从它的channels字段中将此通道提取出来。断言揭示了该通道自身的类型、存储的数据类型和累积模式开关。fromlanggraph.channelsimportLastValue,Topicfromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.typesimportSend,Sequence node(NodeBuilder().subscribe_only(input_channel).do(lambdaargs:args).write_to(output_channel))appPregel(nodes{node:node},channels{input_channel:LastValue(str),output_channel:LastValue(str)},input_channels[input_channel],output_channels[output_channel],)tasks:Topic[Send]app.channels[__pregel_tasks]assertisinstance(tasks,Topic)asserttasks.ValueTypeSequence[Send]asserttasks.accumulateFalse3. 被保护起来的通道虽然__pregel_tasks就是一个普通的Topic类型的通道但是它并未开发对外部使用Pregel把它保护得非常好。我们不能声明一个与之同名的通道否则就会像如下的方式一样抛出一个ValueError并提示Channel __pregel_tasks is reserved and cannot be used in the graph.。fromlanggraph.channelsimportTopicfromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.typesimportSend,Sequencetry:appPregel(nodes{node:NodeBuilder().subscribe_only(__pregel_tasks)},channels{__pregel_tasks:Topic[Sequence[Send]]},input_channels[input_channel],output_channels[output_channel],)assertFalse,Expected an error due to reserved channel nameexceptExceptionase:assertisinstance(e,ValueError)assertstr(e)Channel __pregel_tasks is reserved and cannot be used in the graph.我们也不能采用常规的方式将向其发送Send对象。比如在如下的演示程序中节点foo试图向此通道发送一个驱动节点bar执行的Send对象最终抛出了一个InvalidUpdateError异常并提示Cannot write to the reserved channel TASKS。除此之外由于Pregel在利用它将基于Push模式的任务创建出来后就会将其清空所以我们也无法读取其中的任务。fromlanggraph.channelsimportLastValuefromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.typesimportSendfromlanggraph.errorsimportInvalidUpdateError foo(NodeBuilder().subscribe_to(start,readFalse).do(lambda_:Send(nodebar,argfoobar)).write_to(__pregel_tasks))bar(NodeBuilder().do(lambdaargs:args).write_to(output))appPregel(nodes{foo:foo,bar:bar},channels{start:LastValue(str),output:LastValue(str),},input_channels[start],output_channels[output])try:app.invoke({start:None})assertFalse,Should have raised InvalidUpdateErrorexceptExceptionase:assertisinstance(e,InvalidUpdateError)assertstr(e)Cannot write to the reserved channel TASKS4. 唯一的解决方案我们能够想到的常规方法针对此通道的写入基本都绕不开引擎针对它的保护机制。我们在以Actor模型的角度来看Pregel中提到过节点利用ChannelWriter对象实现针对通道的写入。我们可以将针对通道的写入意图封装成ChannelWriteTupleEntry并以此来创建ChannelWriter这应该是唯一能够欺骗引擎验证的手段。如代码片段所示我们率先执行的节点foo会返回一个驱动节点bar指定的Send对象为了将它写入__pregel_tasks我们创建了一个ChannelWriter针对该通道的写入定义在ChannelWriteTupleEntry对象中具体体现在调用构造函数指定的mapper参数上它提供一个映射将节点的执行结果转成通道名称和值的映射关系。fromlanggraph.pregelimportPregel,NodeBuilderfromlanggraph.channelsimportLastValuefromlanggraph.pregel._readimportPregelNodefromlanggraph.pregel._writeimportChannelWrite,ChannelWriteTupleEntryfromlanggraph.typesimportSend foo:PregelNode(NodeBuilder().subscribe_to(foo).do(lambda_:Send(nodebar,argfoo))).build()entryChannelWriteTupleEntry(mapperlambdaargs:[(__pregel_tasks,args)])foo.writers.append(ChannelWrite(writes[entry]))bar(NodeBuilder().do(lambdaargs:fbar is triggered by{args}.).write_to(output))appPregel(nodes{foo:foo,bar:bar},channels{foo:LastValue(None),output:LastValue(str),},input_channels[foo],output_channels[output],)resultapp.invoke(input{foo:None})assertresult{output:bar is triggered by foo.}

关于恒美微站

恒美微站专注于为个体商户、工作室提供极简自助建站服务,让每个人都能轻松拥有专业网站。

快速链接

  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心

服务项目

  • 可视化建站
  • 拖拽编辑
  • 主题定制
  • SEO 优化
  • 网站托管

联系方式

  • 📍 地址:北京市朝阳区建国路 88 号
  • 📞 电话:400-888-8888
  • ✉️ 邮箱:info@hmyw.cn
  • 🕐 时间:周一至周日 9:00-18:00

© 2024 恒美微站 hmyw.cn 版权所有 | 京 ICP 备 12345678 号