ARTICLE DETAIL

资讯详情

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

differential-dataflow 输入变更全指南:理解 InputSession 的 insert / remove / update / update_at 与时间语义

differential-dataflow 输入变更全指南:理解 InputSession 的 insert / remove / update / update_at 与时间语义 differential-dataflow 输入变更全指南理解 InputSession 的 insert / remove / update / update_at 与时间语义【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway导读本指南以 differential-dataflow 官方教程《Making Changes》章节为核心骨架系统讲解如何通过InputSession向微分数据流differential dataflow计算注入变更包括insert(item)、remove(item)、update(item, diff)、update_at(item, time, diff)四个输入方法的语义差异以及任意 timely dataflow 流都能被重铸为微分集合的互操作机制。读完本文你将掌握交互式驱动微分计算的标准模式——如何在合适的时间点写入变更、推进输入时间并等待输出收敛并能对照仓库中的源码实现external/differential-dataflow/src/input.rs理解其内部批处理与时间放宽原理。该文档位于仓库内 external/differential-dataflow/mdbook/src/chapter_3/chapter_3_3.md属于 differential-dataflow本项目以external/differential-dataflow形式内置官方 mdbook 教程第三章与微分数据流交互Differential Interactions的一部分。一、InputSession从命令式代码到微分数据流的桥在微分数据流中计算一旦被定义与外界的全部交互就归结为两件事改变输入、观察输出。绝大多数的输入由一个InputSession实例管理参见 chapter_3_1// 创建一个输入会话。 let mut input InputSession::new(); // 在一个 dataflow 作用域内把它转成集合。 worker.dataflow(|scope| { // 从我们的输入创建出一个微分集合。 let manages input.to_collection(scope); // ... 在此之上拼接你的计算图 });InputSession扮演桥的角色你在命令式代码里对会话所做的修改会以差分增加/删除的形式进入数据流计算。从源码结构看一个会话内部持有三样东西external/differential-dataflow/src/input.rspub struct InputSessionT: TimestampClone, D: Data, R: Semigroup { time: T, // 会话当前的逻辑时间 buffer: Vec(D, T, R), // 待发出的 (数据, 时间, 差分) 缓冲 handle: HandleT,(D,T,R), // 底层 timely dataflow 输入句柄 }也就是说你调用insert、remove等方法时数据并不会立刻暴露给 timely dataflow而是先进入内部buffer在特定时机以批量batch形式发送到句柄。这正是源码模块注释强调的设计动机InputSession把更新连同其逻辑时间打包并以粗粒度化的 timely dataflow capabilities发送从而向算子实现暴露更多并发度src/input.rs。二、最直接的变更insert 与 remove文档指出InputSession提供了若干方法用来修改底层集合其中最简单的是insert(item)与remove(item)insert(item)在当前输入时间下把该条目的计数加一remove(item)在当前输入时间下把该条目的计数减一。从源码看二者都只是对更通用的update方法做语法糖封装src/input.rsimplT: TimestampClone, D: Data InputSessionT, D, isize { /// 向集合中添加一个元素。 pub fn insert(mut self, element: D) { self.update(element, 1); } /// 从集合中移除一个元素。 pub fn remove(mut self, element: D) { self.update(element, -1); } }微分数据流的哲学是一切输入本质上是多重集合multiset上的差分。插入与删除不区分先后因果只以计数的增减来表达。因此你完全可以先remove一个从未存在过的元素产生负计数再在后续时间点insert它来抵消——系统最终会对齐到正确结果。这一点贯穿教程第三章始终交互式示例中就同时使用了input.remove((person/2, person))与input.insert((person/3, person))见 chapter_3。注意类型参数InputSessionT, D, R中的差分类型R要求实现Semigroup默认R isize。当你的更新需要表达负数与大幅度变更时正是update方法的用武之地。三、任意幅度的变更update(item, diff)update(item, diff)允许你对某个条目指定任意的变化量——可正可负、幅度任意diff属于差分类型R例如isize。其实现src/input.rs/// 增加集合中某个元素的权重。 pub fn update(mut self, element: D, change: R) { if self.buffer.len() self.buffer.capacity() { if self.buffer.len() 0 { self.handle.send_batch(mut self.buffer); } self.buffer.reserve(1024); } self.buffer.push((element, self.time.clone(), change)); }这段代码同时揭示了两个值得注意的实现细节变更被记录在会话当前时间上推入缓冲的三元组是(element, self.time.clone(), change)即时刻取自会话内部时钟你无须也不能在普通update中单独指定时间。缓冲写满即自动发送当buffer达到容量时会先把已有批次通过handle.send_batch发出去再预分配1024个槽位。也就是说即使你从不显式flush累积到一定量的更新也会自动泄洪但对时间推进而言显式调用flush仍是必要的见后文。文档特别指出update的典型适用场景输入数据本身就以差值形式到达或者需要处理大幅度变化时。例如来自外部订阅源的增量流、或需要一次性把某个键的计数调整成千上万的场景直接使用update(item, 1000)远比循环调用insert一百次来得高效清晰。四、乱序到达的输入update_at(item, time, diff)update_at(item, time, diff)让你为一次变更指定任意时间——可以是会话当前时间也可以是当前时间之后的任意未来时刻。方法签名与前置断言src/input.rs/// 在未来某个时间增加元素的权重。 pub fn update_at(mut self, element: D, time: T, change: R) { assert!(self.time.less_equal(time)); // 目标时间必须不早于当前会话时间 if self.buffer.len() self.buffer.capacity() { if self.buffer.len() 0 { self.handle.send_batch(mut self.buffer); } self.buffer.reserve(1024); } self.buffer.push((element, time, change)); }从源码可以看到与update的唯一本质差异是三元组中的时间不再取self.time而是你显式传入的time并伴随一条assert!(self.time.less_equal(time))——你只能向当前及未来写入不能回退到过去。文档给出该方法的动机当数据可能相对时间乱序out of order到达时你仍然希望把它引入计算而不是自己在外层缓冲等待。需要注意引入不代表立即生效你不太可能在输入推进到这些时间之前看到这些变更的效果——但你仍然愿意引入数据而不是自行缓冲它。这正是微分数据流以时间换批处理的核心思想把乱序到达的数据钉在各自的逻辑时间上系统会在推进到相应时间时一次性消化从而让你免于编写复杂的重排序缓冲逻辑。五、Timely streams把任意流重铸为微分集合InputSession并不是输入的唯一来源。文档在 Timely streams 一节明确任何 timely dataflow 流只要记录类型是(data, time, diff)都可以被重新解释为一个微分数据流集合——因此任何能被转成 timely 流的外部变更源都可以充当微分计算的输入。这依赖AsCollectiontrait 提供的as_collection()方法src/collection.rs 中定义并从 src/lib.rs 再导出。chapter_3_1 里给出了典型动机假设你用 timely dataflow 从 Kafka 拉取数据那么在交给微分算子前需要先把 timely 流转成微分集合。除了显式转换仓库还在Inputtraitsrc/input.rs上提供了三种内联创建输入的方式方法语义new_collection()返回(InputSession, Collection)二元组从空集合开始实现见 src/input.rsnew_collection_from(I)接受一个可迭代的初始数据序列每个元素以差分1、最小时间写入src/input.rsnew_collection_from_raw(I)接受(data, time, diff)三元组序列与句柄注入的数据流做concat合并src/input.rs前两种的惯用法是把它放在dataflow闭包内部并把InputSession作为闭包返回值绑定到外部变量见 chapter_3_1// 定义一个新计算。 let mut input worker.dataflow(|scope| { // 从输入创建集合。 let (input, manages) scope.new_collection(); // ... 拼接计算 ... input // 必须把 input 从闭包中返回出来 });这里很容易犯错的点是忘记把 input 返回并绑定。因为dataflow是一次性构建图的作用域如果你在闭包内创建输入却不想办法把它带出来后续就没有句柄可用了。六、让变更真正生效advance_to 与 flush 的配合仅调用insert/remove/update并不会让计算推进。微分数据流在确信拥有足够信息产出正确答案之前只做相对较少的工作——这既需要你提供输入变更同样关键的是你需要承诺停止在某个时间之前继续修改集合见 chapter_3_4。InputSession提供advance_to(time)把会话内部时间推进到time并阻止你在小于该时间的时刻继续提供输入变更源码中以断言强制单调性src/input.rs。这是向微分基础设施作出的强承诺——我在所有不晚于该时间的时刻都不再改这个输入了系统因此可以开始确定相应输出变更。flush()把缓冲的数据强制送入 timely 输入并把底层句柄的 epoch 推进到会话时间src/input.rs。关键点advance_to的调用本身是会被缓冲的直到调用flush()才会向底层 timely 系统暴露。这意味着你可以逐条记录地调用advance_to频率极高而不会淹没底层系统。这正是InputSession相比裸用 timely 输入更安全、更高效的原因之一。一个反复出现的经典错误是没有 advance flush 就期待计算向前推进。这会让程序表现为挂死——事实上计算几乎总是正确的只是因为你扣住了输入在某时刻已停止变化这一信息系统无法判断何时可以产出结果。文档特意以IMPORTANT标注提醒chapter_3_4这也是 differential-dataflow 交互式使用中最容易踩的坑。七、Temporal Concurrency不必每步都 flushflush()并非每次advance_to后都必须调用chapter_3_4 的 Temporal Concurrency 小节advance_to调用只是改变变更发生的逻辑时间你可以把大量这样的变更缓冲起来、只 flush 一次让微分数据流并发地处理这串变更。这通常能显著提升吞吐而只带来名义上的延迟影响。背后的机制是InputSession内部的粗粒度 capability 批处理逻辑时间看起来串行执行但底层 timely dataflow 收到的却是一批带放宽时间限制的更新算子实现得以暴露更多并发。这就是并发出现在逻辑串行数据流中的原因。这一特性在实时流式输入类应用中尤其能体现该主题会在后续章节与示例应用中展开。八、等待计算收敛probe 与 worker.step() 的标准驱动循环变更写入后如何知道输出已经稳定答案是probe。probe()算子在集合上返回一个探针告诉你在哪些时间上该集合仍可能发生变更chapter_3_2。数据流计算不像命令式计算那样强迫某一步执行你只能等待它发生——probe就是用来询问是否已经等到的机制。教程第三章给出完整的交互式循环骨架chapter_3// 变更输入但等待完成。 while person people { input.remove((person/2, person)); input.insert((person/3, person)); input.advance_to(person); input.flush(); while probe.less_than(input.time()) { worker.step(); } person peers; }循环的每一步都是一次标准的写入 → 推进时间 → flush → 跑到输出跟上insert/remove在当前时间记录变更advance_to(person)推进逻辑时间承诺此前不再变化flush()把缓冲与时间承诺暴露给 timelywhile probe.less_than(input.time()) { worker.step(); }驱动 worker 不断步进直到探针确认所有严格早于当前输入时间的变更都已解析完毕。worker.step()会调度每个微分算子执行其积压工作反复调用即可推进系统内全部工作最终使输出 probe 与输入时间对齐。反过来如果在退出dataflow闭包之前你什么都没做也不必担心退出时 timely 会自动反复调用worker.step()直到计算完成输入随Drop自动关闭并 flush所以无事可做的 worker 可以直接退出它仍会参与协作直到所有 worker 完成见 chapter_3_5。显式的step()只在你需要保持对 probe 的交互式访问时才不可或缺。InputSession的Drop实现src/input.rs正是flush()保证了作用域退出即关闭输入——与上述自动完成行为互为印证。九、仓库内的代码证据测试如何驱动输入InputSession的真实行为在集成测试中可以得到直接验证。以 external/differential-dataflow/tests/import.rs 的test_import_stalled_dataflow为例L161-L204let mut input InputSession::new(); let (mut trace, probe1) worker.dataflow(|scope| { let arranged input.to_collection(scope).arrange_by_self(); (arranged.trace, arranged.stream.probe()) }); input.insert(Hello.to_owned()); input.advance_to(1); input.flush(); worker.step_while(|| probe1.less_than(input.time())); input.advance_to(2); input.flush(); worker.step_while(|| probe1.less_than(input.time()));测试清晰展示了本文的核心流程构造InputSession→ 在dataflow内to_collection→ 写入 advance_to→flush→step_while等待 probe 追平输入时间。注意第二次只调用了advance_to(2)而没有写入任何新数据——这再次证明advance_to本身就是一种有效操作它在告诉系统时间 1 已封闭从而让下游这里是后加入的订阅者数据流得以判断自己是否需要等待。同文件中其他测试如test_import_vanilla、L46-L99则展示了与本文Timely streams一节呼应的另一条路径直接在 timelynew_input()句柄上发送(data, time, diff)三元组、逐轮advance_to最后close()——这正是任何 timely 流都可以作为微分输入在测试中的落地形态。而 lib.rs 顶部的库级示例则给出另一个视角advance_to→insert→flush→while probe.less_than(input.time())的驱动循环同样适用于持续喂入、持续观察的场景。十、小结选择正确的输入方法面对一份输入你可以依据其形态做如下决策输入形态推荐方法说明单个元素的增/删insert(item)/remove(item)分别等价于update(item, 1)/update(item, -1)作用于会话当前时间以差值到达 / 大幅度变化update(item, diff)diff可正可负、幅度任意差分类型R需满足Semigroup默认isize相对时间乱序到达update_at(item, time, diff)可在当前及未来的任意逻辑时间落账效果要等输入推进到该时间才会显现其他 timely 数据源as_collection()任何(data, time, diff)类型的 timely 流都可重铸为微分集合配合始终不变的三个纪律advance_to承诺停止过去、flush让承诺可见、probe/step()等待输出收敛——你就能像仓库测试那样以最小心智负担驱动一个高吞吐、低延迟的增量计算。延伸阅读创建输入InputSession 与 new_collection 家族观察输出probe 的语义与用法推进时间与 Temporal Concurrency 的完整讨论worker.step() 与闭包退出的自动完成机制InputSession 完整源码insert/remove/update/update_at/advance_to/flush输入驱动的集成测试实例【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表