原理与源码实现详解)
Pathway 底层引擎揭秘Timely Dataflow 进度追踪Progress Tracking原理与源码实现详解【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway进度追踪Progress Tracking是 Timely Dataflow 运行时体系中最核心也最容易被忽略的基础设施它决定了系统在任意时刻、数据流图中的任意位置能否对未来是否还会收到更小/更多时间戳的数据给出有界承诺从而支撑算子的通知notification机制与整体计算的终止判定。本篇以本仓库external/timely-dataflow子项目 mdbook 第 5.2 节《Progress Tracking》为骨架结合timely/src/progress/与timely/src/dataflow/下的真实 Rust 实现系统讲解其三大组成——数据流图结构的命名、能力Capability的多重度维护、路径摘要Path Summary的传达并最终推导出支撑整个协议的安全性质。读完本文你将能够从数据平面与控制平面分离的角度理解 timely/differential 计算引擎为何能给出确定性的终止与通知保证并能在源码层面定位到每一步机制的具体实现。阅读前提本文在内部机制章节中的定位本文对应文档来自 external/timely-dataflow/mdbook/src/chapter_5/chapter_5_2.md。从仓库结构来看本仓库Pathway的 Rust 引擎构建于 timely dataflow 与 differential dataflow 的分布式增量计算模型之上两个子项目的完整源码被收录在external/目录下其中 mdbook 总目录见 external/timely-dataflow/mdbook/src/SUMMARY.md。mdbook 的第 5 章名为Internals见 external/timely-dataflow/mdbook/src/chapter_5/chapter_5.md它由三个层次递进的子节构成Communication——通信层worker 线程的启动、Configuration、Allocate/Push/Pull等通道抽象Progress Tracking本文主题——进度追踪层在数据流图上回答还会不会有数据来Containers——容器抽象数据在数据流边上传输时所用存储表示与核心算子解耦。chapter_5_1中特别提醒读者本节是 internals 部分你也可以直接对着通信 crate 写代码但 timely dataflow 的一个优点就是你不需要这么做——换句话说理解进度追踪并非日常写算子的前提但它决定了上层算子的能力边界。进度追踪要解决什么问题先看一个朴素的设定系统中有多个 worker它们共同推动数据流图上的数据前进。数据可以在 worker 之间流动而算子处理数据的后果几乎无法预先承诺一个 worker 收到一条记录可能什么都不做也可能发出上千条输出记录其中一些最终会流向我们。在这种前提下每个位置需要能够对我是否还可能收到更多数据做出有意义meaningful的陈述。Timely Dataflow 的答案是给数据附加一个逻辑时间戳logical timestamp它不是物理时间不是记录被创建时 worker 的挂钟时间而可以是任何满足若干约束的类型最常见的例子是从 0 开始递增的序号sequence number只要数据流图结构满足一组很自然的约束系统就能在图中每个位置做出有界陈述形式为你只会看到大于等于这些时间的时间戳当未来时间戳的集合为空时系统甚至能得知该数据流已经完成completed——这正是及时timely通知机制得以成立的基石。Timely Dataflow 对计算结构的约束是算子在某个时间戳上发送带时间戳消息必须先持有该时间戳上的能力capability。由此进度追踪可以被拆成两件事(i) 集体维护视图多个 worker 共同维护数据流图每个位置上、每个时间戳还剩下多少未释放的能力outstanding capability(ii) 独立传达影响每个 worker 独立地确定自己视图中的能力变化对其他位置意味着什么并把这种影响广播出去。在进入这两个方面之前必须先能命名数据流图中的各个部分。第一步给数据流图上的每个位置命名算子、输入端口与输出端口数据流图承载若干算子operator。对进度追踪而言算子只按索引index识别不需要了解算子内部逻辑。每个算子有若干个输入端口input port和输出端口output port算子之间通过把一个输入端口连接到某个输出端口通常是另一个算子的相连。一个输出端口可以连接到多个不同的输入端口——这意味着在输出端口产生的一条消息会被投递给所有挂接的输入端口fan-out广播。从**进度协调器progress coordinator**的视角看一个算子的输出端口是带时间戳数据的source源一个算子的输入端口是带时间戳数据的target目标。因此 timely 用两个不同类型Source与Target来标识端口用类型系统避免把输入输出搞混。mdbook 给出的示意定义如下pub struct Source { /// Index of the source operator. pub index: usize, /// Number of the output port from the operator. pub port: usize, } pub struct Target { /// Index of the target operator. pub index: usize, /// Number of the input port to the operator. pub port: usize, }仓库中的真实定义需要指出的是mdbook 中这段代码是简化示意。在当前仓库的源码中Source/Target的真实定义位于 external/timely-dataflow/timely/src/progress/mod.rs字段名从index演进为node并引入了统一的Location类型与Port枚举见 progress/mod.rs 的Port::Target(usize)/Port::Source(usize)例如/// Names a source of a data stream. #[derive(Copy, Clone, Eq, PartialEq, Ord, PartialOrd, Debug)] pub struct Source { /// Index of the source operator. pub node: usize, /// Number of the output port from the operator. pub port: usize, } /// Names a target of a data stream. #[derive(Copy, Clone, Eq, PartialEq, Ord, PartialOrd, Debug)] pub struct Target { /// Index of the target operator. pub node: usize, /// Number of the input port to the operator. pub port: usize, }配合Location::new_source/Location::new_target构造函数progress/mod.rsSource、Target都能统一投影为Location { node, port }供下游可达性reachability计算复用同一套数据结构。至此数据流图的结构可由一张连接表完整描述Vec(Source, Target)。从这个列表中既能推断出算子总数、每个算子的输入/输出端口数目也能枚举出所有连接。于是图中每个位置都有了名字——非Source即Target可以画成每个算子一个圆圈、每个端口一条短线、输出端口连向目标输入端口的图。第二步维护能力Capability的多重度目标当算子在运行、消息在被收发时worker 们集体追踪系统中每个时间戳、每个位置上尚未释放的能力数量。能力存在于两个地方算子可以显式持有能力以便在它的每个输出上以某个时间戳发送消息每条带时间戳的消息本身就携带一个对应时间戳的能力。进度追踪记录的是能力的多重度multiplicity在位置l上、时间t的能力到底有多少个。对绝大多数位置和时间这个数是 0只要计算尚未完成某些位置和时间上这个数必然为正而由于各 worker 上报的变更可能乱序到达数值还可能短暂地为负。初始化与三类动作计算启动时不存在在途消息每个算子一开始都拥有在其任一输出上发送任意时间戳消息的能力。由于这是所有 worker 的公共知识每个 worker 把每个算子输出的初始计数初始化为#workersworker 总数——每个 worker 都明白必须听到#workers份该能力已被丢弃的上报才能真正认为该能力已从系统中消失。计算推进过程中算子只能执行三类动作文档中列得很明确动作类别对能力计数的影响1. 消费consume输入消息获取消息携带的能力对应Target输入端口上的能力被递减同时算子Source输出端口上的能力被递增2. 克隆clone、降级downgrade或丢弃drop其持有的能力改变算子对应Source输出端口上的计数3. 在其持有能力的任意时间戳上发送输出消息使该Source所连接的各Target端口上的计数被递增消息被消费能力从Target迁移到算子内部并可能重新出现在Source发送消息则把能力推向所有下游Target。具体地一批batch变更具有如下形状(Vec(Source, Time, i64), Vec(Target, Time, i64))即每个位置、每个时间上的增减量。批的完整性至关重要进度追踪协议的安全性依赖于绝不传播半成形的进度消息——例如消费了一条输入消息却忘记上报其能力的获得就会破坏协议。广播与本地视图每个 worker 把进度变更批的流通过点对点FIFO 通道广播给系统中所有 worker包括它自己。任意时刻每个 worker 都只是看过了其他每个 worker 产生的进度批序列的某个前缀且随着时间推进每个 worker 最终能看到每个序列的任意前缀。于是任意时刻某个 worker 对系统能力的视图 收到的所有进度变更批的累加 每个输出上的初始能力初始多重度为 worker 数。这个视图可能又惊又乱甚至出现负计数但它始终满足一个与能力含义的传达相关的安全条件——这正是下一节要展开的内容。源码对照ChangeBatchTprogress/change_batch.rs正是(T, i64)变更累积器的实现追加更新update、惰性合并compact时按 key 排序并把累加归零的项剔除见 change_batch.rs、drain_into空时直接mem::swap优化传播路径。CapabilityTdataflow/operators/capability.rs内部持有time: T与RcRefCellChangeBatchT——Capability::new会立刻对内部的ChangeBatch执行update(time, 1)capability.rs。这就是克隆/延迟一个能力 在对应的ChangeBatch里 1这一抽象在实现层面的落点能力在算子输出计数上的增减就由这些共享的ChangeBatch汇聚而成。SharedProgressTprogress/operate.rs把父子 scope 间交换的进度信息分成了四类通道frontiers外部前沿变化、consumeds已消费消息、internals内部能力变化、produceds已产出消息——这与文档描述的消费输入、操作内部能力、发送输出三类动作一一对应。第三步传达含义Communicating Implications与路径摘要仅仅知道能力在哪还不够。即使时间t上存在能力也不代表时间t可能到达图中任意位置。要做出精确陈述必须讨论能力沿数据流图路径传播时对时间戳的约束即路径摘要Path Summary。时间戳与路径摘要的类型契约进度追踪运行在一个所有算子共享同一时间戳类型T: Timestamp的数据流图上。Timestamptrait 要求PartialOrder——即两个时间戳可以有序也可以不可比incomparable。mdbook 中的简化定义是pub trait Timestamp: PartialOrder { type Summary : PathSummarySelf; }仓库中 progress/timestamp.rs 的真实定义更完整叠加了CloneEqHashOrdSendAnyDataDebug等约束并增加minimum()其核心契约不变每个实现Timestamp的类型必须指定一个关联类型Summary它实现PathSummarySelf。PathSummary非正式地概括时间戳沿 timely dataflow 中某条路径传播时必然会发生什么。大多数你能画出的路径只有平凡摘要保证不改变但有些路径会强制改变时间戳——例如从循环loop尾部回到循环头部的路径必然使经过的任意时间戳能力上的循环计数器坐标 1。pub trait PathSummaryT : PartialOrder { fn results_in(self, src: T) - OptionT; fn followed_by(self, other: Self) - OptionSelf; }两个方法各司其职results_in(self, src: T) - OptionT——解释时间戳沿这条路径走必须变成什么。注意可能返回None时间戳可能根本无法沿该路径前进。例如某个路径摘要path会把时间戳加一那么path.results_in(4) Some(5); path.results_in(u64::max_value()) None; // 溢出走不过去results_in只允许前进时间戳这是硬性要求对所有摘要p若p.results_in(x)得到Some(time)则必有x.less_equal(time)。followed_by(self, other: Self) - OptionSelf——解释两条路径摘要如何复合。构造摘要时从单条边对应的路径出发再通过复合若干路径及其摘要拼出更复杂的路径。和results_in一样两条路径复合可能产生永远不会让时间戳通过的结论此时可返回None例如两条各把循环计数器增加其最大值一半的路径。两条路径摘要的序关系定义为对任意时间戳两个摘要作用于该时间戳得到的结果同样有序。因为路径摘要是偏序的当要概括从位置 A 到位置 B、沿众多路径之一时很快会面对路径摘要的集合——多条路径对应多个摘要可以丢弃那些严格大于集合中其他元素的摘要但仍可能留下多个不可互相比较的摘要。因此传达含义的第二步本质上就是为图中任意两个位置Source或Target之间确定并锁定最小的路径摘要集合。这就是数据流图的编译后表示compiled representation从它才能推出图中某处的消息可能导致另一处的消息之类的陈述。源码佐证纯加法时间戳的摘要实现在 progress/timestamp.rs 中所有具备checked_add的整型通过宏一次性实现Timestamp/PathSummaryresults_in即self.checked_add(*src)followed_by即self.checked_add(*other)——溢出天然返回None与文档中u64::max_value()的例子完全吻合单元类型()与Duration也各有实现。而保留不可比的最小元素集合这一操作对应的正是 progress/frontier.rs 中的AntichainTinsert时丢弃所有less_equal于新元素的旧元素始终维护最小反链frontier.rs。算子自身的摘要Operatetrait 与get_internal_summary路径摘要从何而来答案是每个算子必须实现Operatetrait并为自身做摘要对从每个输入到每个输出的内部路径提供一组路径摘要描述一条带时间戳的消息到达某个输入在某个输出上可能产生什么样的时间戳。fn get_internal_summary(mut self) - (VecVecAntichainT::Summary, RcRefCellSharedProgressT);对应仓库实现见 progress/operate.rs其中输入输出间以(input, output)对为索引每个条目是一个AntichainT::Summary父 scope 借此获知某算子是否在输出上初始持有能力。对于大多数算子这个摘要是平凡的我任一输入上带时间戳的数据都可能以相同时间戳出现在我的任一输出上。它虽然对推进没多大帮助却也是该算子能给出的唯一保证——默认实现正是如此。Feedback算子则是绝大多数有趣路径摘要的起点它是反馈环feedback loop中的算子保证时间戳的某个特定坐标被递增。timely dataflow 中的所有环都必须经过此类算子这也保证了下文无不做时间戳推进的环的约束成立于是环状数据流中会出现从下游算子输出到上游算子输入的非平凡摘要。查看 dataflow/operators/feedback.rs 可看到运行时对摘要的使用当Feedback将消息从循环末端接回循环头时执行summary.results_in(cap.time())并把时间戳推进后的消息放回环首feedback.rs构造阶段则用feedback_core(summary)记录这个推进摘要。更抽象的loop_variable(summary)则把循环次数编码进Product时间戳。另一种同样有用的摘要是摘要的缺席有时候图中两个点之间根本没有路径。例如在输入数据上做了些计算后把数据送进一个迭代子计算那么从迭代回到输入计算的路径不存在——迭代子计算中在途的记录不会阻塞在输入数据上的计算因为我们知道没有从迭代返回输入的通路。更精妙的例子来自带数据输入 诊断输入的算子如果算子内部摘要揭示诊断输入不可能导致数据输出那么就可以随时投递诊断查询作为该输入的数据而不阻塞数据输出的下游消费者——timely 能看到尽管有消息在途它们到不了数据输出因此不必处在数据计算的关键路径上。编译表示Compiled Representation两两之间的摘要集合从算子摘要构造出路径摘要再从路径摘要为每一对Source/Target求出摘要集合之后问题一处带时间戳的消息如何导致另一处的消息就有了答案。整个过程唯一要求的全局约束是不允许存在不严格推进时间戳的环no cycles that do not strictly advance a timestamp——这正是反链上偏序论证成立的前提。仓库中这段逻辑的落点在 progress/reachability.rs。其模块注释开门见山管理 timely dataflow 图内的 pointstamp 可达性timely dataflow 关心的是能力沿有向图的边和节点传播后可能到达哪里。 文件内文档示例展示了BuilderTracker的完整用法先add_node(index, inputs, outputs, internal_summaries)声明每个算子的内部摘要再用add_edge(Source::new(...), Target::new(...))铺设连接例子中刻意构成 0→1→2→0 的环随后update_source(Source, time, delta)引入 pointstamp 变化、propagate_all()完成传播从pushed()读出传播到各Target的结果——注意示例中从node 2绕回node 0的Target上时间 17 被推进到了 18reachability.rs恰好演示了环路路径对时间戳的强制递增。安全性质Safety Property现在把两部分拼起来每个算子维护图中每个位置Source或Target的时间戳能力计数同时静态地定义从任意位置到任意其他位置的路径摘要集合。安全性质可陈述为对由参与 worker 任意进度更新前缀累加所得的任意计数集合如果对图中某个位置l1和时间戳t1而言不存在另一个位置l2与时间戳t2满足 (a)(l2, t2)处累加计数严格为正且 (b) 存在从l2到l1的路径摘要p使得p(t2) t1那么就不会有任何消息带着时间戳t1到达位置l1。文档提及该性质有 Tom Rodeheffer 提供的 TLA 形式化即 Naiad 时钟协议规范与模型检验/正确性证明。这里不展开外部资料仅在直觉层面复述其论证。为什么在几乎不协调的情况下仍成立先建立一些反直觉这个性质成立但到目前为止对 worker 间消息通信几乎没做任何假设。证明只依赖消息投递是at most once至多一次——消息不应在途被复制。它不关心消息投递顺序不要求消息一定被投递这只是安全性而非活性甚至允许算子消费、消化、发送那些本地进度协议还不知道其存在的消息。数据平面消息发送与控制平面进度批发送之间几乎不存在协调唯一要求是你不得为某个尚未执行的动作发送进度更新批。那么正确性的直觉从哪里来三个不变式与冻结论证尽管只看到其他 worker 进度批的前缀我们仍能推理每个 worker 未来的进度批必须长什么样。核心是三条与动作因果一致的性质任何被消费的消息必然对应一条已被产出的消息——即使我们还没听说它任何被产出的消息必然涉及一个被持有的能力——即使我们还没听说它任何被持有的能力必然源于另一个被持有的能力或一条被消费的消息。接下来对(位置 li, 时间戳 ti)定义次序——这就是 Naiad 论文中称为pointstamp的对象。pointstamp 按could-result-in可能导致关系成偏序(li, ti)could-result-in(lj, tj)当存在从li到lj的路径摘要p使p(ti) tj。这是偏序因为(i) 自反性由定义保证(ii) 反对称性——两个不同的 pointstamp 若互相could-result-in会推出一个不严格推进时间戳的环而这被约束排除(iii) 传递性由路径摘要构造的正确性保证。每条原子进度更新批还有如下性质为便于陈述可假设算子不能克隆能力、必须消费一个能力才能发消息每批更新先递减某个 pointstamp然后可选地递增严格大于它的若干 pointstamp等价地任何递增了某 pointstamp 的进度批必然同时递减了某个严格小于它的 pointstamp。归纳可得任何净效果为递增某 pointstamp的进度批集合必然伴随净效果为递减某个严格小于它的 pointstamp。最后是冻结freeze论证设想在某 worker 可能收到(l1, t1)消息的瞬间把系统冻结——不再执行新动作只结算已执行动作对应的全部进度更新让这些进度批充分流通使该 worker 得到冻结系统的完整视图此时所有 pointstamp 计数都应非负。冻结期间新增流入的进度批只有同时对某个严格更小的 pointstamp 施加净负效果才能在(l1, t1)上产生净正效果而由于最终累加非负这只有当该更小 pointstamp 在冻结前已有正累加时才可能发生。若冻结前不存在任何这样的正累加 pointstamp稳定化后(l1, t1)就不可能为正因而也就不可能存在一条等待被接收的消息——安全性质得证。概念—源码对照速查表文档概念第 5.2 节仓库实现位置Source/Target端口命名timely/src/progress/mod.rs外加统一抽象Location/Portmod.rs能力多重度变更累积(T, i64)ChangeBatchTprogress/change_batch.rs算子持有的能力对象CapabilityTdataflow/operators/capability.rsnew时对内部ChangeBatch加一delayed/downgrade改变时间戳Timestamp/PathSummary契约progress/timestamp.rs整型经checked_add宏实现timestamp.rs最小反链丢弃严格更大的摘要AntichainTprogress/frontier.rs算子必须自报摘要Operate::get_internal_summaryprogress/operate.rs与SharedProgress四类通道operate.rs唯一会推进时间戳的环算子Feedback/ConnectLoop/loop_variabledataflow/operators/feedback.rs运行时以summary.results_in推进时间戳见 feedback.rs两两位置间的编译表示最小摘要集合 传播pointstamp 可达性Builder/Trackerprogress/reachability.rs小结从这一节可以提炼出一条理解 timely 系引擎的思维主线数据平面只管发消息控制平面只管报能力变化能力是数据与时序之间的唯一粘合剂路径摘要是把局部能力信息翻译成全局时序承诺的编译器而安全性质保证了两平面之间这种极弱耦合下结论仍然正确。differential-dataflow本仓库 external/differential-dataflow之所以能实现持续增量地维护结果并对外输出确定性的完成信号正是因为它站在本节这套进度追踪协议之上——你可以顺着Operatetrait 与reachability模块继续深挖它如何驱动算子的通知回调也可以回到 mdbook 阅读 Communication 与 Containers 两个姊妹章节补全 Internals 的完整图景。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考