ARTICLE DETAIL

资讯详情

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

PocketFlow 节点通信实战:基于 Shared Store 与 Params 构建共享状态流水线

PocketFlow 节点通信实战:基于 Shared Store 与 Params 构建共享状态流水线 人工智能大模型AI Agent工作流自动化RAG【免费下载链接】PocketFlowPocket Flow: 100-line LLM framework. Let Agents build Agents!项目地址https://gitcode.com/gh_mirrors/poc/PocketFlow点击查看免费下载本篇技术指南以 cookbook/pocketflow-communication/README.md 中的通信示例为骨架深入讲解 PocketFlow 中节点间通信的两种核心机制——**Shared Store共享存储**与Params参数并结合仓库源码剖析其底层实现。读完本文你将掌握如何用共享字典在多个节点之间传递与维护状态、如何设计条件流转conditional transition以及如何通过 Params 在批量任务中传递任务标识符。一、节点间通信的两种方式Shared Store 与 Params在 PocketFlow 中节点Node与流水线Flow之间的通信只有两种方式官方文档 docs/core_abstraction/communication.md 对此有明确阐述Shared Store共享存储——适用于绝大多数场景一个全局数据结构通常是内存字典所有节点都可以通过prep(shared)读取、通过post(shared, ...)写入非常适合承载数据处理结果、大段内容或任何需要被多个节点访问的数据设计共享存储的结构并提前填充数据是使用它的前提。Params参数——仅用于 Batch 批量场景每个节点拥有一个由父 Flow 传入的局部、临时的params字典用作任务的标识符参数的键和值应当是**不可变immutable**的适合承载文件名、数字 ID 这类标识性数据。文档给出的一个经典类比非常适合理解这两种机制Shared Store 好比堆heap被所有函数调用共享Params 好比栈stack由调用方在调用时分配。一个面向全局共享、生命周期贯穿整个流水线一个面向单次任务、随调用而更新。最佳实践来自官方文档绝大多数场景都应使用 Shared Store以将数据 Schema 与计算逻辑分离Separation of Concerns——先设计共享数据结构再让各节点围绕该结构读写。Params 更多是 Batch 场景下的语法糖详见 docs/core_abstraction/batch.md。二、示例项目总览一个跨节点共享状态的词频统计器本仓库的 cookbook/pocketflow-communication/ 目录实现了一个完整的单词计数器示例用来演示节点如何通过共享存储通信。项目结构如下pocketflow-communication/ ├── README.md ├── requirements.txt ├── main.py ├── flow.py └── nodes.py其中main.py程序入口创建 Flow 并运行flow.py定义节点间的连接关系转移图nodes.py三个业务节点 一个终止节点的实现requirements.txt依赖声明核心依赖为pocketflow0.1.0。安装与运行pip install -r requirements.txtpython main.py程序启动后会提示输入文本其行为是统计输入文本中的单词数将统计结果写入共享存储展示运行以来的累计统计信息已处理文本数、单词总数、平均值。输入q即可退出程序。三、源码级剖析三个节点如何通过 Shared Store 协作示例共定义了三个业务节点外加一个终止节点它们通过共享字典shared完成数据的写入与读取。下面逐节点讲解。1. TextInput读取输入并初始化共享存储TextInput 节点 的核心逻辑class TextInput(Node): def prep(self, shared): # 获取用户输入 return input(Enter text (or q to quit): ) def post(self, shared, prep_res, exec_res): if prep_res q: return exit # 返回 exit 动作走向终止节点 shared[text] prep_res # 把文本写入共享存储 if stats not in shared: # 首次运行时初始化统计字段 shared[stats] { total_texts: 0, total_words: 0 } shared[stats][total_texts] 1 # 更新累计文本数 return count # 返回 count 动作走向 WordCounter可以看到 Shared Store 在这里扮演了数据 Schema 定义者的角色shared[text]存放原始文本shared[stats]存放累计统计。这两个 key 构成了整个流水线共享的数据契约——后续节点都围绕它们工作。这正是官方文档强调的先设计共享存储结构再编写节点逻辑。2. WordCounter从共享存储读取并更新统计WordCounter 节点 展示了完整的prep - exec - post三段式class WordCounter(Node): def prep(self, shared): # 从共享存储读取上游写入的文本 return shared[text] def exec(self, text): # 纯计算统计单词数不触碰共享存储 return len(text.split()) def post(self, shared, prep_res, exec_res): # 将单词数累加到共享存储的统计字段 shared[stats][total_words] exec_res return show # 返回 show 动作走向 ShowStats这里体现了 PocketFlow 节点的一个关键设计exec是纯计算函数只依赖prep的结果不直接接触shared所有对共享存储的读写都收敛在prep读与post写两个阶段。从 pocketflow/init.py 的_run实现可以验证这一调用链def _run(self, shared): p self.prep(shared) # 读共享存储 e self._exec(p) # 纯计算 return self.post(shared, p, e) # 写共享存储这种分层让计算逻辑易于测试与复用——exec不依赖任何外部状态传入同样的输入必然得到同样的输出。3. ShowStats读取并展示累计统计ShowStats 节点 负责把统计结果呈现给用户class ShowStats(Node): def prep(self, shared): # 从共享存储读取统计字段 return shared[stats] def post(self, shared, prep_res, exec_res): stats prep_res print(f\nStatistics:) print(f- Texts processed: {stats[total_texts]}) print(f- Total words: {stats[total_words]}) print(f- Average words per text: {stats[total_words] / stats[total_texts]:.1f}\n) return continue # 返回 continue回到 TextInput 形成循环注意shared[stats]是TextInput在post阶段写、WordCounter在prep阶段读、又在post阶段更新、最后由ShowStats读取的——同一个字典在多个节点间流转并被持续修改这正是 Shared Store 模式维护跨多次执行的持久状态的核心价值每次循环迭代total_texts与total_words都在累计而无需任何全局变量或外部存储。四、Flow 中的连接语法默认转移与条件转移节点单独存在没有意义是 flow.py 中的连接定义了它们的执行顺序def create_flow(): text_input TextInput() word_counter WordCounter() show_stats ShowStats() end_node EndNode() # 条件转移根据 post 返回的动作字符串选择下一节点 text_input - count word_counter word_counter - show show_stats show_stats - continue text_input # 形成循环 text_input - exit end_node # 退出分支 return Flow(starttext_input)这里用到了两种连接语法它们的底层实现都可以在 pocketflow/init.py 中找到默认转移等价于next(node, actiondefault)当节点的post返回None时走这条边- action 条件转移__sub__会构造一个_ConditionalTransition对象再通过其__rshift__调用src.next(tgt, action)把动作字符串 - 下一节点的映射注册进successors字典。Flow 运行时get_next_nodepocketflow/init.py会以节点post返回的动作字符串为键查找后继节点若找不到对应动作且节点存在其他后继会抛出UserWarning: Flow ends: ...提示流水线终止。运行循环状态在迭代之间如何延续main.py 的启动代码非常简洁flow create_flow() shared {} # 创建一个空字典作为 Shared Store flow.run(shared)Flow.run内部通过_orchpocketflow/init.py驱动节点循环执行def _orch(self, shared, paramsNone): curr, p, last_action copy.copy(self.start_node), (params or {**self.params}), None while curr: curr.set_params(p) last_action curr._run(shared) curr copy.copy(self.get_next_node(curr, last_action)) return last_action关键点在于同一个shared字典对象被传入每一次节点执行所以TextInput在第 2 轮迭代写下的stats第 3 轮依然存在——这就是维护跨多次节点执行的状态的实现原理。同时节点对象使用copy.copy复制避免并发或嵌套场景下的状态污染。循环退出条件是当前节点没有匹配的后继当用户输入q时TextInput.post返回exit指向EndNode一个空的Node子类其post隐式返回None由于EndNode没有注册任何后继流水线自然终止。五、Params面向 Batch 的任务标识符通道除了 Shared StorePocketFlow 还提供 Params 作为第二种通信手段。官方文档强调Params不可变、通过set_params()设置、并且每次被父 Flow 调用时会被清除并更新。一个典型用法如下摘自 docs/core_abstraction/communication.mdclass SummarizeFile(Node): def prep(self, shared): # 通过 self.params 读取任务标识符 filename self.params[filename] return shared[data].get(filename, ) def exec(self, prep_res): prompt fSummarize: {prep_res} return call_llm(prompt) def post(self, shared, prep_res, exec_res): filename self.params[filename] shared[summary][filename] exec_res # 按标识符写回共享存储 return default # 单节点测试直接给节点设置参数 node SummarizeFile() node.set_params({filename: doc1.txt}) node.run(shared) # Flow 场景Flow 的 params 会覆盖节点 params flow Flow(startnode) flow.set_params({filename: doc2.txt}) flow.run(shared) # 实际处理的是 doc2.txtParams 与 Shared Store 的分工非常清晰维度Shared StoreParams生命周期贯穿整个 Flow 运行单次节点调用被父 Flow 覆盖可变性可变各节点可读写不可变键值不可修改典型用途数据、大内容、共享状态文件名、ID 等任务标识符适用场景几乎所有场景主要用于 Batch 批量处理官方文档还给出了一条重要警告只需要设置最上层 Flow 的 params因为下层节点的 params 会被父 Flow 覆盖若需设置子节点参数应参考 docs/core_abstraction/batch.md 中关于 Batch 的说明。六、从源码与测试看通信机制的正确用法仓库的单元测试 tests/test_flow_basic.py 从多个角度印证了 Shared Store 与转移机制的约定其中最有价值的两条经验可直接迁移到你的流水线设计中post返回None走默认转移返回字符串走条件转移。测试test_sequence_with_rshift中NumberNode、AddNode、MultiplyNode的post隐式返回None因此用串联即可线性执行而CheckPositiveNode必须显式返回positive或negative才能触发分支。转移找不到对应动作时会触发警告而非静默崩溃。test_flow_ends_warning_default_missing与test_flow_ends_warning_specific_missing验证了Flow ends: ... not found in [...]的UserWarning帮助你在开发期尽早发现连接错误。七、Shared Store 使用的最佳实践清单综合官方文档 docs/core_abstraction/communication.md 与本示例源码在使用 Shared Store 时建议遵循以下实践先设计数据结构再写节点根据应用需求确定共享字典的 Schema或数据库表结构可以包含data、summary、config等分组示例中stats字段的初始化即遵循这一原则读写边界清晰prep负责读、post负责写、exec保持纯计算确保计算逻辑可独立测试用动作字符串驱动流转让post返回有语义的动作名如count、exit、continue配合- action 构建可读的转移图用循环节点维护累计状态Shared Store 的跨迭代持久性天然支持计数器、汇总器等有状态逻辑无需引入外部状态管理Params 只放不可变标识批量处理时用 Params 携带文件名、ID用 Shared Store 承载真正的数据与结果。八、深入阅读通信机制官方详解docs/core_abstraction/communication.md本示例源码cookbook/pocketflow-communication/nodes.py、cookbook/pocketflow-communication/flow.py、cookbook/pocketflow-communication/main.py框架核心实现Node/Flow/转移机制pocketflow/init.py转移与分支的测试佐证tests/test_flow_basic.py与 Params 密切相关的批量机制docs/core_abstraction/batch.md相关示例批量场景下的通信实践可参考 cookbook/pocketflow-batch/README.md 与 cookbook/pocketflow-map-reduce/README.md通过 Shared Store 与 Params 这套堆 栈的通信模型PocketFlow 让任意复杂的多节点流水线都能以极简的字典契约完成数据交换与状态维护——先定义数据 Schema再用条件转移编织执行路径最后交给 Flow 的循环驱动这正是构建可维护、可扩展 LLM 应用的起点。赞分享人工智能大模型AI Agent工作流自动化RAG【免费下载链接】PocketFlowPocket Flow: 100-line LLM framework. Let Agents build Agents!项目地址https://gitcode.com/gh_mirrors/poc/PocketFlow点击查看免费下载相关推荐PocketFlow 节点间通信指南Shared Store 与 Params 的职责划分与源码级解析PocketFlow 节点间通信指南Shared Store 与 Params 的职责划分与源码级解析 Shared Store 与 Params 是 Poc人工智能大模型AI Agent工作流自动化RAGPocketFlow 共享状态Shared State实战指南用 shared 字典打通 Flow 中所有 Node 的数据流转PocketFlow 共享状态Shared State实战指南用 shared 字典打通 Flow 中所有 Node 的数据流转 shared 字典是 P人工智能AI 应用AI AgentTheatre微前端通信模式基于URL的状态共享Theatre微前端通信模式基于URL的状态共享 你是否在开发微前端应用时遇到过跨应用状态同步难题多个独立部署的前端应用如何高效共享用户状态本文将介绍Th前端上一篇WebRTC发展史从RTP到现代实时通信的技术演进下一篇TVM 测试指南Python 单元测试的 Target 参数化、本地运行与 CI 集成实战创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表