)
differential-dataflow 迭代计算指南从iterate不动点算子到Variable通用迭代Pathway 增量引擎的底层基石【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本文以本仓库external/下 vendored 的 differential-dataflow 开源 crate 的 mdbook 教程第 1.3 节为基础原文chapter_1_3.md系统讲解该库最富特色的迭代计算能力先用iterate算子写出反复执行直到不动点的单层递归如管理者—员工传递闭包再深入更通用的手动迭代方式——通过scope.scopedenter/leaveVariable构造支持多输入、多循环变量、多输出的迭代子图。本仓库根目录 Cargo.toml 声明differential-dataflow { path ./external/differential-dataflow }即这套迭代机制正是 Pathway 增量计算引擎依赖的底层能力。读完本文你将掌握用iterate一行式求解图算法不动点、用Variable实现相互递归以及何时必须手动consolidate以避免差异无限循环。迭代differential dataflow 最有趣的能力之一对实时增量计算框架而言迭代iteration是难点也是亮点普通批处理框架可以简单地重跑而差分数据流需要在输入持续变化的前提下把循环内部每一轮产生的**差异differences**增量地传导下去直到这些差异消散即逻辑上达到不动点fixed point。在 differential-dataflow 中迭代通过iterate算子完成它接收一个输入集合以及一段描述如何把集合变换一轮的逻辑闭包随后把这段逻辑对输入集合无限次地应用输出最终收敛的结果。用源码的话说iterate的实现并不是真的把闭包跑无穷多次而是建立了一个迭代的 timely dataflow 子计算让差异在其中循环流动直到差异消散表示计算到达不动点或直到某轮迭代停止The implementation ofiteratedoes not directly apply the closure, but rather establishes an iterative timely dataflow subcomputation, in which differences circulate until they dissipate (indicating that the computation has reached fixed point), or until some number of iterations have passed.见 operators/iterate.rs 的模块文档在底层这依赖 timely dataflow 的 **scope作用域**机制迭代是一个嵌套的数据流子图它把外层时间戳附加一个迭代轮次维度——从iterate的签名可以看到闭包操作在CollectionIterativea, G, u64, D, R之上operators/iterate.rs其中u64正是记录循环轮数的迭代时间分量。下面用一个例子说明iterate的用法。用iterate计算管理者—员工传递闭包假设我们有一个(Manager, Employee)集合表示谁管理谁我们希望产出并持续维护每位管理者名下含间接下属的员工总数。最自然的思路是从管理者—员工关系出发反复扩展出传递的管理关系——这正是iterate的用武之地。原文给出了如下代码chapter_1_3.mdmanager_employee .iterate(|manages| { // if x manages y, and y manages z, then x manages z (transitively). manages .map(|(x, y)| (y, x)) .join(manages) .map(|(y, x, z)| (x, z)) });这段代码把每一轮迭代拆解为三步.map(|(x, y)| (y, x))把(x, y)x 管理 y翻转为(y, x)以便拿下属作为连接键.join(manages)与当前manages集合做连接——若 y 管理 z同时 (y, x) 翻转后表示 x 管理 y则产出三元组(y, x, z)语义为x 管理 y且 y 管理 z.map(|(y, x, z)| (x, z))投影出(x, z)即新增的传递管理关系x间接管理 z。每一轮产出比上一轮更长的管理链当所有传递关系都被推导出来、新集合与旧集合不再有任何差异时计算达到不动点并终止。join、map的增量语义可参考教程前文 chapter_1_2.md算子综述。iterate的底层实现源码解读iterate并非黑魔法它其实是在手动迭代基础上封装出的语法糖。其核心实现位于 operators/iterate.rsimplG: Scope, D: OrdDataDebug, R: Abelian IterateG, D, R for CollectionG, D, R { fn iterateF(self, logic: F) - CollectionG, D, R where G::Timestamp: Lattice, fora F: FnOnce(CollectionIterativea, G, u64, D, R)-CollectionIterativea, G, u64, D, R { self.inner.scope().scoped(Iterate, |subgraph| { let variable Variable::new_from(self.enter(subgraph), Product::new(Default::default(), 1)); let result logic(variable); variable.set(result); result.leave() }) } }可以看到iterate内部只做了四件事与后续将要介绍的通用迭代完全一致scoped(Iterate, ...)在顶层 scope 中打开一个名为Iterate的子作用域Variable::new_from(self.enter(subgraph), Product::new(Default::default(), 1))把输入集合enter进子作用域作为循环变量的初值时间戳步长取Product::new(Default::default(), 1)表示外层时间不变、迭代轮次每轮 1let result logic(variable); variable.set(result);把用户提供的变换逻辑作用到变量上再用set把下一轮的值反馈回去形成循环result.leave()把收敛后的结果带出子作用域返回给外层。Iteratetrait 与Variable等定义都在同一文件 operators/iterate.rs。实现中还特别注释道直接返回variable包装的集合会产生显著更多的差异记录而result经过合并post-consolidation因此能大幅减少从循环中产出的记录量。重要提醒iterate不会自动consolidate模块文档明确警告operators/iterate.rsThe dataflow assembled byiteratedoes not automatically insertconsolidatefor you. This means that either (i) you should insert one yourself, (ii) you should be certain that all paths from the input to the output of the loop involve consolidation, or (iii) you should be worried that logically cancelable differences may circulate indefinitely.即iterate拼装的数据流不会自动插入consolidate。在差分数据流里记录携带代数差如1/-1若某条路径上逻辑上可抵消的差异迟迟不合并循环就可能永不终止。因此你需要自己收尾consolidate()或者确保循环内部所有路径都经过了本身会合并的算子——reduce、distinct、count都属于安全收尾算子。这一点在Iteratetrait 的文档示例中也能看到operators/iterate.rsvalues.map(|x| if x % 2 0 { x/2 } else { x }) .consolidate()更通用的迭代手动构造迭代上下文iterate的写法简洁但表达能力有限。原文指出当你的迭代计算涉及(i) 多个输入、(ii) 多个循环变量、(iii) 多个输出时就需要手动构造迭代上下文。第一步用scope.scoped打开子作用域迭代必须发生在一个嵌套子数据流里在那里时间戳可以被附加上额外信息例如循环轮次。原文给出// if you dont otherwise have the scope .. let scope manager_employee.scope(); scope.scoped(|subscope| { // More stuff will go here });scoped由 timely dataflow 提供iterate内部也正是调用了self.inner.scope().scoped(Iterate, ...)。一旦进入subscope外层manager_employee就无法直接使用了——因为每个集合都隶属于某个特定 scope。第二步用enter把外部集合带进子作用域要使用子作用域之外的集合需要把它带入子作用域这就是enter算子// if you dont otherwise have the scope .. let scope manager_employee.scope(); scope.scoped(|subscope| { // we can now use m_e in this scope. let m_e manager_employee.enter(subscope); });enter把外层的外层时间戳的数据提升到子作用域内、为时间戳补上迭代维度相应地后面还有与它配对的leave负责把结果带回外层。用Variable声明可更新的循环集合接下来需要定义一个能在每轮迭代中被更新的变量。differential-dataflow 提供了Variable结构体递归定义的集合recursively defined collection先指定它的初值一个集合再通过set给出下一轮如何更新的定义。Variable实现了Deref到其内部Collectionoperators/iterate.rs因此在绝大多数可以写集合的地方都能直接使用它。原文示例// if you dont otherwise have the scope .. let scope manager_employee.scope(); scope.scoped(|subscope| { // we can now use m_e in this scope. let m_e manager_employee.enter(subscope); let variable Variable::from(m_e); let step variable .map(|(x, y)| (y, x)) .join(variable) .map(|(y, x, z)| (x, z)); variable.set(step); });这里step复用上一节的传递闭包逻辑并且同时引用了variable两次翻转后再join自身这正是递归地基于自身定义自身的体现。版本注记mdbook 教程写作年代使用Variable::from(m_e)而当前仓库中 vendored 版本的构造器已演变为 Variable::new_from(collection, step) 与 Variable::new(scope, step)后者用于初值为空的场景能生成更简单的数据流图。因此与当前源码一致的可运行写法是let step Product::new(Default::default(), 1); // 外层时间不变迭代轮次 1 let variable Variable::new_from(m_e, step);关于Variable与set源码 operators/iterate.rs 还有几个值得注意的实现事实set(self, result: Collection...)按值消费self并返回一个Collection。模块文档明确指出这防止了同一变量被设置多次的误用operators/iterate.rsset内部会把source初值集合做negate()再与结果concat以撤消对初值的重复计算如果希望保留初值集合继续参与循环应使用 set_concat它避免了先追加初值、再撤回初值的多余数据流工作等价于用Variable::new建空变量再set(self.concat(result))在set_concat中每一轮的结果会通过step.results_in(t)把时间戳推进到下一迭代轮再经connect_loop从feedback句柄送回循环起点。Variable要求差异类型实现Abelian支持取负。对于差异类型只实现Semigroup、只能单向增长如只能 1、无法抵消的场景同一文件还提供了 SemigroupVariable对应地Iteratetrait 也为G: Scope本身而非Collection提供了基于SemigroupVariable的实现operators/iterate.rs。leave把结果带回外层作用域完成定义之后通常需要返回变量收敛后的最终值。与进入作用域用的enter配对leave负责产出其所调用集合的最终值。原文的完整示例// if you dont otherwise have the scope .. let scope manager_employee.scope(); let result scope.scoped(|subscope| { // we can now use m_e in this scope. let m_e manager_employee.enter(subscope); let variable Variable::from(m_e); let step variable .map(|(x, y)| (y, x)) .join(variable) .map(|(y, x, z)| (x, z)); variable .set(step) .leave() });可以看到set返回集合后紧接着被.leave()整个表达式的类型就回到了外层 scope 的Collection。原文评价说虽然写法啰嗦了一些但这应当与前面用iterate方法表达的其实是同一个计算——只不过当你需要更多输入、更多输出或多个相互递归的变量时这套手动框架随时可以支撑你扩展。仓库内的真实用例iterate驱动的图算法要验证上述机制的真实性与威力直接看本仓库中 vendored crate 自带算法的用法即可——它们把iterate与enter结合产出了工业级的增量图算法广度优先距离标注 BFSalgorithms/graphs/bfs.rs// initialize roots as reaching themselves at distance 0 let nodes roots.map(|x| (x, 0)); // repeatedly update minimal distances each node can be reached from each root nodes.iterate(|inner| { let edges edges.enter(inner.scope()); let nodes nodes.enter(inner.scope()); inner.join_core(edges, |_k,l,d| Some((d.clone(), l1))) .concat(nodes) .reduce(|_, s, t| t.push((s[0].0.clone(), 1))) })每一轮迭代把已标记节点的距离 1 后沿边向外扩散join_core再与既有标记concatreduce取到每个节点的最小距离当传播不再产生更短距离的差异时迭代自然收敛。注意edges、nodes这两个外层输入都通过enter进到循环内部使用而循环变量inner本身承载着每轮更新的距离标注。类似地本 crate 中还提供了若干可直接对照的迭代算法实现强连通分量与传递闭包的反复剪枝algorithms/graphs/scc.rs、algorithms/graphs/scc.rs前缀和prefix sum的倍增合并algorithms/prefix_sum.rs中以 .iterate(|ranges| ...) 迭代扩大区间标签传播等状态传递算法algorithms/graphs/sequential.rs、algorithms/graphs/propagate.rs。阅读这些实现时你会发现一个反复出现的模式循环体外部的输入一律用enter(inner.scope())带入循环体内尽量以reduce收尾——前者是迭代上下文的常规动作后者则是为了利用reduce内部合并机制规避前文未consolidate导致差异循环不散的风险。迭代机制在 Pathway 仓库中的定位需要强调的是本文讨论的 differential-dataflow 是被 vendored 到本仓库external/目录下的独立 crate另有配套的 timely-dataflow 同样位于 external/timely-dataflow。仓库根目录 Cargo.toml 中differential-dataflow { path ./external/differential-dataflow }说明 Pathway 的 Rust 引擎直接以这份 vendored 代码为编译依赖因此本节讲解的iterate/Variable迭代机制、时间戳分层与差异合并约定是理解引擎内部增量语义时不可绕过的一环。若要在仓库中做更深度的源码研读推荐路线是先通读本教程章节 chapter_1_3.md本文主体配合算子系统章节 chapter_2/chapter_2_7.md 了解iterate在算子全景中的位置再精读实现文件 operators/iterate.rs约 267 行是本节全部概念的最终权威最后以 bfs.rs、scc.rs、prefix_sum.rs 等算法实现作为阅读测试验证自己对迭代语义的理解。小结iterate与手动Variable如何取舍需求推荐写法单一输入集合递归逻辑可复用如传递闭包、BFScollection.iterate(\|x\| {...})需要多个输入 / 多个相互递归的循环变量 / 多个输出scope.scopedenterVariable可多个leave循环变量初值为空Variable::new(scope, step)数据流图更简单循环变量有初值来自外部集合Variable::new_from(source, step)差异类型不支持取负、循环只会单向增长SemigroupVariable/ 对Scope调用iterate循环内部存在可能互相抵消的差异路径收尾consolidate()或确保经由reduce/distinct/count无论走哪条路有两条纪律始终不变凡进入循环的外部集合都要enter凡要取回结果的集合都要leave凡可能出现可抵消差异的路径都要以会合并的算子收尾。掌握了它们你便可以在增量数据流上自由地写出各种不动点算法——从谁管理谁的传递闭包到全图的 BFS、标签传播与强连通分量求解。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考