十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

Differential Dataflow 何时使用:从 Pathway 底层引擎看函数式、数据并行、迭代与增量更新的能力边界

Differential Dataflow 何时使用:从 Pathway 底层引擎看函数式、数据并行、迭代与增量更新的能力边界 Differential Dataflow 何时使用从 Pathway 底层引擎看函数式、数据并行、迭代与增量更新的能力边界【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathwayDifferential Dataflow 是一种有主见opinionated的流式计算框架它面向集合上的组合算法刻意优先服务一部分任务、并对另一部分任务毫不掩饰地相对无用。本文以本仓库内external/differential-dataflow随书文档《When to use Differential Dataflow》为骨架结合其配套源码external/differential-dataflow/src/下的核心实现与示例逐条拆解该框架的四大设计支柱——函数式编程、数据并行、迭代与增量更新帮助读者建立正确的预期判断什么样的问题值得用它。本仓库PathwayPython 流处理 / 实时分析 / LLM Pipeline / RAG 的 ETL 框架在external/目录中完整托管了 timely-dataflow 与 differential-dataflow 两个底层依赖其中 differential-dataflow 的 fork 还带有src/pathway定制模块理解本文内容即理解这套实时计算引擎擅长什么、不擅长什么的设计原点。该篇文档的来龙去脉原文位置external/differential-dataflow/mdbook/src/chapter_0/positives.md属于随书教程mdbook/中 Motivation动机 一章chapter_0.md的导引部分。其姊妹篇 negatives.md 则讨论何时不要使用 Differential Dataflow两者共同构成完整的预期管理。面向的读者原文档明确指出要真正爱上 differential dataflow关键是把预期设置正确set expectations appropriately为此必须先理解它到底想做好什么。它不承诺一切皆可而是列出几个它能稳定兑现的核心能力。一句话概括它的设计目标面向集合collections上的组合算法典型例子是大规模图计算——例如计算并持续维护一张图的连通分量connected components。从 SQL、MapReduce 这类经典大数据计算一直延伸覆盖到演绎推理系统deductive reasoning与部分结构化机器学习形态这套技术都适用。在原文档introduction.md中这种工作方式的描述极简先写一个程序然后改变它的输入系统会在最短毫秒级时间内把输出中对应的变化反馈给你。支柱一函数式编程Functional Programming含义算子不修改输入只把输入变换为输出Differential dataflow 的算子全部是函数式的map、filter、join、reduce等操作把输入集合转换为输出集合而不去就地改动输入。函数式算子乃至整个函数式程序的最大价值在于当输入变化时我们更容易推演程序整体会经历怎样的变化。从源码看这一承诺落实在Collection这一核心抽象上。collection.rs 的文档写道编写 differential dataflow 计算时你仿佛是在对一个静态数据集施加函数式变换、不断派生出新集合一旦计算写定你就可以通过插入与删除元素来改变集合differential dataflow 会把变化沿你的函数式计算传播出去并把对应的输出变化报告给你。具体实现中每个Collection内部只是一条携带(D, G::Timestamp, R)三元组的 timely 数据流——数据、时间戳、变化量collection.rs。所有算子都是对这条流做纯函数变换例如map仅对每个(data, time, delta)应用映射逻辑后再封装回集合collection.rsfilter仅保留满足谓词的三元组collection.rsconcat则对应两个集合的相加collection.rs。输入的集合始终原样保留派生集合完全由变换产生。为什么函数式约束值得付出认知切换的成本原文档坦承函数式编程有约束性但强调交换到的回报很强大高效的分布式执行、迭代执行与增量执行——这正是接下来三个支柱的基石。因为变换可组合、且不破坏输入系统才能安全地只重算真正受影响的部分而不是把整个计算推倒重来。支柱二数据并行Data Parallelism含义不相交的数据块可以独立运算Differential dataflow 的算子大多是数据并行的同一算子可以独立地作用在输入中互不相交的部分上。这带来两个层面的收益跨 worker 并行与众多大数据框架一样可以把工作分布到多个 worker 上执行。这是由底层 timely dataflow 自动完成的——lib.rs 明确说明 differential dataflow 构建于 timely dataflow 之上后者自动地把计算并行到多线程、多进程乃至多台计算机上。约束更新流向这一点更为关键。数据并行意味着可以按记录/键把更新的传播范围切分开从而只为发生变化的值执行重算re-computation而不是让一次输入抖动触发全图刷新。在可执行的示例 hello.rs 中可以看到 worker 模型的具体形态每个 worker 通过worker.index()与worker.peers()获知自己的编号与总并行度随后只处理属于自己的那份输入图中数据按index轮转分配。这印证了多 worker 各自处理不相交输入块的编程模型。从实现层面看真正的按需重算由数据交换与算子的分组如reduce按 key 分组见 reduce.rs配合实现key 被哈希分发到特定 worker因此单个 key 的变化只会唤醒维护该 key 的那条执行路径。支柱三迭代Iteration含义计算并持续维护迭代式程序原文档特别强调这是与以往工作相比最重要的技术分水岭绝大多数数据库与大数据处理器都不具备对迭代计算的高效支持与持续维护能力而差分数据流允许你在计算内部嵌套迭代直到结果收敛不动点。源码中的iterate算子是这一能力的直接体现。iterate.rs 的文档说明iterate接受一个把集合映射为同类型集合的闭包其输出是把该闭包无限次应用的结果。实现上它并非直接循环而是建立一个迭代式 timely dataflow 子计算差异differences在回路中持续循环直到它们互相抵消表示已达不动点或达到指定迭代轮数。iterate的底层是更灵活的Variable递归定义的集合它允许你自建更复杂的循环形态如两个集合协同演化、旋转循环以返回中间结果文档见 iterate.rs。仓库中的迭代型算法与示例都直接使用了这套机制例如图算法库 src/algorithms/graphsbfs.rs、bijkstra.rs、scc.rs、propagate.rs等演示程序 examples/bfs.rs广度优先搜索、examples/pagerank.rs、examples/monoid-bfs.rs。一个容易踩坑的工程细节值得注意iterate.rs 的文档专门提醒iterate不会自动为你插入consolidate合并/压缩操作符。因此你必须自己插入一次consolidate或者确保从循环输入到输出的每条路径都经过了合并类算子reduce、distinct、count本身都会做合并否则逻辑上可以抵消的差异可能会无限循环导致程序不终止。这是把迭代支柱用好时必须补上的纪律。支柱四增量更新Incremental Updates含义输入一变输出随之更新且代价正比于变化Differential dataflow 会在输入发生变化时维护计算结果而这种维护代价通常远低于从零开始完整重算。原文档特别强调它是被专门设计为同时提供高吞吐与低延迟的——两者兼顾而非二选一。这一能力的数据模型基础是带符号多重集 可抵消差异每条记录关联一个差值difference最常见的是isize表示出现次数的增减difference.rs 定义了Semigroup半群加法 判零其中is_zero被系统用来判断某条更新累加为零、可以安全删除difference.rs 进一步定义Abelian阿贝尔群含取负iterate循环回路的收敛正是依赖这种可取负抵消的性质通过input.update(record, 1)/input.update(record, -1)即可表达插入/删除这正对应用户文档中最经典的两步操作写程序 → 改变输入。以 lib.rs 的库级文档示例来看该例也是一个完整的度数分布计算程序形态为// 在一个 worker 上构建数据流 let (mut input, probe) worker.dataflow(|scope| { // 创建边集合输入做两轮计数 let (input, edges) scope.new_collection(); // 抽取源点字段然后计数 let degrs edges.map(|(src, _dst)| src) .count(); // 抽取计数字段再计一次数得到度分布 let distr degrs.map(|(_src, cnt)| cnt) .count(); // 观察输出变化并设置探针 let probe distr.inspect(|x| println!(observed: {:?}, x)) .probe(); (input, probe) }); // 驱动计算推进时间、插入数据、冲刷并步进直到探针就绪 loop { let time input.epoch(); for round in time .. time 100 { input.advance_to(round); input.insert((round % 13, round % 7)); } input.flush(); while probe.less_than(input.time()) { worker.step(); } }这里体现的增量特性值得展开计数类算子如count在输入记录到来时不重扫全量数据而只对被影响到的 key 重算当某个节点度数从 d 变成 d1 时输出侧通常只产生 4 条变化记录——新旧度数下各一条加与减lib.rs。这正是维护代价正比于变化本身的直观例证。可运行的同类示例位于 examples/hello.rs它比文档片段更进一步先随机加载一张含nodes个顶点、edges条边的图随后按batch为粒度持续地向输入中同时注入 1新增边与 -1删除边并用probe判定每轮输出是否已收敛input.advance_to(round); input.update((rng1.gen_range(0, nodes), rng1.gen_range(0, nodes)), 1); input.update((rng2.gen_range(0, nodes), rng2.gen_range(0, nodes)), -1);也就是说差分引擎维护的是边集持续增删变化下的度分布结果每次只在变化传播完成后报告差异——这是增量更新在真实批处理循环hello.rs、examples/degrees.rs中的典型使用姿势。四条支柱如何协同一个判断清单原文档把这些能力称为目前其他解决方案中尚不具备的功能组合。如果你面临的问题需要、或哪怕是部分受益于以下特征differential dataflow 就值得研究支柱你获得的能力判断信号函数式编程输入变换可推演、可组合改动可追溯你能用map / filter / join / reduce描述业务逻辑数据并行跨 worker 扩展 只重算变化的部分数据量远超单机、且关心局部变化导致的局部重算迭代维护不动点计算、支持非平凡控制流需要 BFS/PageRank/连通分量/传播类算法且结果要随输入实时更新增量更新高吞吐与低延迟并存地维护结果输入持续变化、输出要长期保持在最新配套视角也请读完何时不该用为了让预期设定得足够准确原文档在同一章配了 negatives.md要点同样值得写进决策清单差分重算的工作量与计算路径的实际变化量成正比。即使你觉得新旧结果定性上差不多通往结果的路径可能已经面目全非此时差分引擎除了在输入变化处重放计算外别无选择。历史追踪可能造成显著内存占用。为了维护结果框架需要记录计算如何演化构造出历史远大于任一时刻状态的计算并不困难这类场景下可能出现出乎意料的大内存足迹。把positives.md与negatives.md放在一起读才能形成何时使用/何时不用的完整决策框架。结语能力边界会随时间扩大原文档最后指出能被 differential dataflow 良好容纳的问题范围只会越来越大——随着社区在这些方向上持续积累算法与算子实现。本仓库即是活证据在 vendored 的 external/differential-dataflow 中除了examples/bfs、pagerank、graspan、stackoverflow、degrees等、tpchlike/22 条 TPC-H 风格查询与src/algorithms/graphs/等不断扩充的算法集合还能看到面向交互式查询的interactive/子项目以及src/pathway这一 fork 定制模块——后者表明该差分引擎在本仓库Pathway的技术栈中承担了实时的底层计算角色external/timely-dataflow与external/differential-dataflow并列存在于external/目录可相互对照阅读。因此对何时使用的回答可以归结为一句可执行的话当你的问题可以被表达为对集合的函数式变换、需要在多 worker 上并行、需要迭代到收敛、并且希望结果随输入增量式地保持新鲜时——Differential Dataflow 就是值得优先验证的方案反之若你的计算只是一次性批处理、输入固定不变或计算历史极度膨胀、内存难以承受则应回头参考其能力边界文档再作取舍。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表