ARTICLE DETAIL

资讯详情

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

Rayon 并行迭代器内部机制深度解析:plumbing 模块的 Pull/Push 双模式与 ProducerCallback 回调设计

Rayon 并行迭代器内部机制深度解析:plumbing 模块的 Pull/Push 双模式与 ProducerCallback 回调设计 【免费下载链接】rayonRayon: A data parallelism library for Rust项目地址https://gitcode.com/gh_mirrors/ra/rayon点击查看免费下载本篇文章以 Rayon 仓库中 src/iter/plumbing/README.md 为核心系统讲解并行迭代器底层的plumbing管道设计它如何通过Producer/Consumer两套 trait 实现可分裂的并行迭代bridge如何把两种模式衔接起来以及ProducerCallback这一被生命周期问题逼出来的回调机制为何存在。读完本文你将理解drive_unindexed、split_at、split_off_left等内部方法的真实含义并掌握为 Rayon 编写自定义组合子所需的底层知识。本文是内部设计文档的讲解不涉及并行迭代器的日常使用方式那是 src/iter/mod.rs 的职责。为什么并行迭代器比串行迭代器复杂普通的串行迭代器Iterator只需要回答一个问题下一个元素是什么 因此它的核心方法只有next()。而并行迭代器ParallelIterator必须多回答一个问题如何把自己分裂成两半让两半在不同的线程上并行工作正是可分裂这个额外能力让并行迭代器的设计远比串行迭代器复杂它不仅要能产出数据还要能确定从哪里分裂、分裂后每一半各有多少数据。Rayon 的当前设计为此提供了两种截然不同的运行模式且并非所有迭代器都同时支持这两种模式——这正是两套 trait 并存的原因Pull 模式拉取模式对应Producer与UnindexedProducertrait。迭代器被要求用next产出下一个元素本质上和普通迭代器一样但多了可以分裂成两半在独立线程上产出不相交的元素的能力。Push 模式推送模式对应Consumer与UnindexedConsumertrait。方向完全反过来迭代器把每个元素依次交给Consumer去处理类似一次for_each调用——每当一个新元素产出就会调用consume方法处理它。实际 trait 比这更复杂因为要支持贯穿整个处理过程、最终被归约的状态。下面分别展开这两种模式。Pull 模式Producer 与 UnindexedProducerProducer在源码中的定义位于 src/iter/plumbing/mod.rs#L56文档注释把它概括为可分裂的IntoIterator一个Producer可以在任意时刻被转换成一个迭代器之后按需产出元素但在转换之前它可以先用split_at在某个指定位置分裂成两个 producer——一个产出该位置之前的元素另一个产出之后的元素两者可以继续独立分裂或各自转成迭代器。在 Rayon 中这种分裂正是用于把工作划分到不同线程。pub trait Producer: Send Sized { type Item; type IntoIter: IteratorItem Self::Item DoubleEndedIterator ExactSizeIterator; fn into_iter(self) - Self::IntoIter; fn min_len(self) - usize { 1 } // 默认分裂到单个元素 fn max_len(self) - usize { usize::MAX } // 默认允许完全不分裂 fn split_at(self, index: usize) - (Self, Self); fn fold_withF(self, folder: F) - F where F: FolderSelf::Item, { folder.consume_iter(self.into_iter()) } }关键点在于split_at接受一个明确的索引把元素切成0..index与index..N两段。因此只有索引型indexed迭代器才能以 Pull 模式工作——它们确切知道自己会产生多少数据也知道如何定位到指定索引。而UnindexedProducer定义于 src/iter/plumbing/mod.rs#L231面向那些不知道精确长度、或长度无法用usize表示的情形例如String的字符或 32 位平台上长度可能超过usize的Rangeu64。它的分裂方法是split——不求分裂到精确位置只要求近似对半pub trait UnindexedProducer: Send Sized { type Item; // 若可能则从中间分裂出一个新 producer否则返回 None fn split(self) - (Self, OptionSelf); fn fold_withF(self, folder: F) - F where F: FolderSelf::Item; }注意split的返回类型是(Self, OptionSelf)数据量不够时允许只返回一个None表示无法继续分裂。理论上任何Producer都可以当作UnindexedProducer使用此时split只需实现为split_at(length/2)但当前仓库并没有利用这种可能性——当精确长度已知时split_at本身已经足够。Push 模式Consumer 与 UnindexedConsumerConsumer定义于 src/iter/plumbing/mod.rs#L123在文档注释中被描述为广义的 fold 操作——实际上每个 consumer 最终都会被转换成一个Folder。与Producer对称Consumer也可以用split_at分裂但分裂会多产出一个reducer分裂出的两个 consumer 被独立喂入元素完成后由 reducer 把两个结果合并成一个。pub trait ConsumerItem: Send Sized { type Folder: FolderItem, Result Self::Result; type Reducer: ReducerSelf::Result; type Result: Send; // 分裂为两个 consumer一个处理 0..index另一个处理 index.. // 同时产出一个 reducer用于最终归约两个结果 fn split_at(self, index: usize) - (Self, Self, Self::Reducer); fn into_folder(self) - Self::Folder; fn full(self) - bool; // 是否希望停止处理更多元素如搜索已完成 }与之配套的是Foldersrc/iter/plumbing/mod.rs#L154与Reducersrc/iter/plumbing/mod.rs#L197FolderItem封装标准的 fold 操作用consume(item)逐个喂入元素并返回新的顺序状态全部消费完后用complete()产出最终值full()则给出是否已满、可以提前停止的提示。它还有一个可选覆写的consume_iter默认实现就是循环调用consume并检查full覆写它可以获得更高效的专用实现。ReducerResult是Consumer的收尾步骤consumer 分裂成两半、各自处理完毕后得到两个结果由reducer.reduce(left, right)合并为一个。与 Producer 一样Consumer 也有无索引变体UnindexedConsumersrc/iter/plumbing/mod.rs#L208但它的分裂不是split_at而是split_off_leftpub trait UnindexedConsumerI: ConsumerI { fn split_off_left(self) - Self; fn to_reducer(self) - Self::Reducer; }split_off_left没有任何索引参数分裂出的两个半部分必须准备好处理任意数量的数据并且不知道这些数据在整个数据流中的位置。文档特别强调了一个微妙的语义split_off_left返回的 consumer 处理左半部分数据且顺序具有意义——对于find_first这类方法返回值产生的数据优先于self产生的数据。最终两半都消费完毕后用to_reducer得到的 reducer 合并结果。与 Producer 侧不同并非所有 Consumer 都能工作在无索引模式。for_each和reduce可以但collect_into_vec不行——因为收集到目标集合时每个元素的位置至关重要。这一限制也印证了文档中的一句感叹有趣的是只有部分 consumer 能在 unindexed 模式下工作但所有producer 都能驱动一个 unindexed consumer反过来只有部分 producer 能驱动 indexed consumer但所有consumer 都能接收索引。这种方差正是整套设计精妙之处。一条完整迭代器链的执行过程文档用一个相对复杂的迭代器链来演示几乎全部可能发生的情况vec1.par_iter() .zip(vec2.par_iter()) .flat_map(some_function) .for_each(some_other_function)从链尾反向构建 Consumer处理一条迭代器链时从尾部开始创建 consumer。这条链的最终步骤是for_each所以它首先创建一个ForEachConsumer——拿到一个元素就调用some_other_function。源码见 src/iter/for_each.rs#L15这是一个非常简单的 consumer因为它在元素之间不需要传递任何状态full()永远返回falsereducer 是 src/iter/noop.rs 中的NoopReducer。然后for_each把这个 consumer 传给链上的前一个迭代器flat_map途径是调用ParallelIteratortrait 上的drive_unindexed方法定义见 src/iter/mod.rs#L2410。drive_unindexed的语义就是让这个迭代器产出元素并把它们逐个喂给这个 consumer——它只对 unindexed consumer 有效。FlatMap 为何只能驱动无索引 ConsumerFlatMap恰好只支持 unindexed consumer见 src/iter/flat_map.rs#L37 的实现。原因很本质flat-map 根本不知道它会产出多少元素。如果你要求 flat-map 直接产出第 22 个元素它做不到——至少在没有中间状态的情况下做不到。它不知道处理第一个输入会产生 1 个、3 个还是 100 个输出因此要产出任意位置的元素它基本上只能从头开始顺序执行这显然不是我们想要的。但对于 unindexed consumer 这完全无所谓因为它们不需要知道自己会收到多少数据。于是FlatMap用FlatMapConsumer把ForEachConsumer包起来。这个FlatMapConsumer每次收到一个元素就调用some_function得到一个并行迭代器然后让这个新迭代器去驱动drive内部的ForEachConsumer。具体逻辑见 src/iter/flat_map.rs#L107 的FlatMapFolder::consume它用map_op(item).into_par_iter()生成子迭代器并drive_unindexed再用to_reducer得到的 reducer 把每次结果合并到previous中。FlatMap的drive_unindexed随后把FlatMapConsumer继续往链的上游传递传给前一个迭代器——zip。到这里有趣的事情发生了。Zip 引发的模式切换从 Push 到 PullZip本质上无法被实现为一个 consumer——至少在没有中间线程、channel或协程的情况下做不到。问题在于它必须**锁步lockstep**地同时走两个迭代器它没法同时调用两个drive方法只能一次调用一个。所以在这一点上zip迭代器需要从Push 模式切换到 Pull 模式。这也解释了为什么Zip只在其输入实现IndexedParallelIterator时可用——从任意位置开始产出数据的能力正是它需要的。如果要在位置 22 处分裂一个 zip 迭代器就必须能从索引 22 开始直接 zip而不必从索引 0 重新走一遍。因此Zip的drive_unindexed不再继续创建 consumer见 src/iter/zip.rs#L30而是创建了一个producer——ZipProducer——并调用internals模块中的bridge函数。创建ZipProducer会为被 zip 的两个迭代器各创建一个 producer这之所以可行正是因为它们都实现了IndexedParallelIterator。ZipProducer的split_atsrc/iter/zip.rs#L138同步分裂左右两个 producerlen()取两个输入长度的较小值src/iter/zip.rs#L54。bridge连接 Producer 与 Consumer 的枢纽bridge函数src/iter/plumbing/mod.rs#L346负责把 consumer 链此处是flat_mapfor_each与 producer 链zip及其上游连接起来pub fn bridgeI, C(par_iter: I, consumer: C) - C::Result where I: IndexedParallelIterator, C: ConsumerI::Item, { let len par_iter.len(); return par_iter.with_producer(Callback { len, consumer }); // ...Callback 内部最终调用 bridge_producer_consumer(self.len, producer, self.consumer) }它不断把 producer/consumer 对分裂下去直到块足够小然后从 producer 拉取元素喂给 consumer。核心递归逻辑在bridge_producer_consumersrc/iter/plumbing/mod.rs#L385的helper函数中fn helperP, C(len, migrated, mut splitter, producer, consumer) - C::Result { if consumer.full() { consumer.into_folder().complete() } else if splitter.try_split(len, migrated) { let mid len / 2; let (left_producer, right_producer) producer.split_at(mid); let (left_consumer, right_consumer, reducer) consumer.split_at(mid); let (left_result, right_result) join_context( |context| helper(mid, context.migrated(), splitter, left_producer, left_consumer), |context| helper(len - mid, context.migrated(), splitter, right_producer, right_consumer), ); reducer.reduce(left_result, right_result) } else { producer.fold_with(consumer.into_folder()).complete() } }三个分支分别对应consumer 已满提前终止还能继续分裂对半分裂后用join_context并行执行两半最后 reduce不能再分裂转入顺序 fold。这里的migrated参数来自join_context的context.migrated()用于告知 splitter 任务是否被偷到了其他线程从而影响分裂策略。基础情形索引型与无索引型的终点bridge还有另一个典型使用场景当链在某个索引型 producer如 slice 或 range处触底时。与此对应无索引 producer如字符串字符使用bridge_unindexedsrc/iter/plumbing/mod.rs#L438它基于Splitter而非LengthSplitter并调用UnindexedProducer::split与UnindexedConsumer::split_off_left。分裂策略Splitter 与 LengthSplitter作为对文档的源码级补充bridge的分裂策略并非盲目对半切而是由 src/iter/plumbing/mod.rs#L251 的Splitter与LengthSplitter两个结构体控制。Splitter实现的是偷窃式自适应策略thief-splitting初始时把期望的分裂次数设为当前线程数crate::current_num_threads()每次成功分裂就减半但一旦发现任务被偷到其他线程执行stolen true就把剩余期望分裂次数重置回线程数与当前剩余值取较大者因为被偷说明负载不均衡、值得再细分。LengthSplittersrc/iter/plumbing/mod.rs#L289在偷窃策略之上再考虑剩余长度min来自producer.min_len()默认 1可由with_min_len调大max来自producer.max_len()默认usize::MAX可由with_max_len调小。构造时用len / max估算出达到上限所需的最少分裂次数try_split则要求len / 2 min才允许继续分裂——保证永远不会切出比min更小的块。文档中特别提醒Rayon 通常会自适应调整分裂粒度以降低开销因此with_min_len/with_max_len一般并不需要主动使用。ProducerCallback被生命周期逼出来的回调机制理想方案闭包回调前面提到调用reduce()这类并行动作方法时会创建归约型consumer然后调用drive_unindexed()或drive()一路向上游创建更多 consumer。但最终总会到达链的起点或到达像zip()这样需要协调多个输入的迭代器——此时就必须开始把并行迭代器转换成 producer。转换的入口是IndexedParallelIterator上的with_producer()方法它采用一种回调callback方案。在理想世界中它应该可以这样用闭包实现base_iter.with_producer(|base_producer| { // 这里base_producer 就是 base_iter 的 producer });这样像map()这样的组合子就可以先获取基础迭代器的 producer把它包装成自己的MapProducer再传给回调struct MapProducerf, P, F: f { base: P, map_op: f F, } implI, F IndexedParallelIterator for MapI, F where I: IndexedParallelIterator, F: MapOpI::Item, { fn with_producerCB(self, callback: CB) - CB::Output { let map_op self.map_op; self.base_iter.with_producer(|base_producer| { let map_producer MapProducer { base: base_producer, map_op: map_op }; callback(map_producer) }); } });这个示例已经展现出回调方案的强大之处我们可以拿走par_iter的所有权再把它的各个部分的所有权交给 producer例如当迭代器拥有一个mutslice 时这非常有用或者创建共享引用放进 producer。以 map 为例并行迭代器拥有map_op我们借用它的引用放进MapProducer——这意味着MapProducer可以轻易分裂自己并共享这些引用。此外with_producer还能创建并行执行期间所需的资源因为 producer 不必被返回。绊脚石关联类型中的生命周期不幸的是有一个障碍实际无法像上面那样使用闭包。原因在于map_producer的类型。如果用闭包来写with_producertrait 会不得不长成这样pub trait IndexedParallelIterator: ParallelIterator { type Producer; fn with_producerCB, R(self, callback: CB) - R where CB: FnOnce(Self::Producer) - R; ... }注意这里不得不引入关联类型Producer才能把回调的参数声明为Self::Producer。而一旦尝试为Map写这个 impl问题就暴露了implI, F IndexedParallelIterator for MapI, F where I: IndexedParallelIterator, F: MapOpI::Item, { type MapProducer MapProducerf, P::Producer, F; // ^^ 等等这个 f 是什么 fn with_producerCB, R(self, callback: CB) - R where CB: FnOnce(Self::Producer) - R { let map_op self.map_op; // ^^^^^^ f概念上就是这个引用的生命周期 // 因此每次调用 with_producer 时它都不一样 } }这看起来似曾相识这正是试图定义Iterabletrait 时遇到的同一个问题。producer 类型需要包含一个生命周期f它指向with_producer的函数体内部因此在 impl 层面不可见、不在作用域内。如果 Rust 有associated type constructors即 RFC 1598 提议的特性就能用那种方式解决——但另一种更立即可行的解决方案是用一个专用的回调 traitProducerCallback取代FnOnce。解决之道面向 Producer 泛化的回调 traitProducerCallback定义于 src/iter/plumbing/mod.rs#L17文档注释把它称为一种泛化的闭包类似FnOncepub trait ProducerCallbackT { type Output; fn callbackP(self, producer: P) - Self::Output where P: ProducerItemT; }采用这个 trait 后with_producer()的签名变为fn with_producerCB: ProducerCallbackSelf::Item(self, callback: CB) - CB::Output;注意这个签名从头到尾都不需要说出 producer 的具体类型——不再有Producer关联类型了。这是因为callback()方法对所有 producerP都是泛化的。代价是||闭包语法糖失效了必须手动创建回调结构体这确实有点繁琐。于是Map的实际代码变成了这样与 src/iter/map.rs#L70 的真实实现几乎逐字一致implI, F IndexedParallelIterator for MapI, F where I: IndexedParallelIterator, F: MapOpI::Item, { fn with_producerCB(self, callback: CB) - CB::Output where CB: ProducerCallbackSelf::Item { return self.base.with_producer(Callback { callback: callback, map_op: self.map_op }); // ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ // 闭包语法糖的手动版本创建一个实现 ProducerCallback 的结构体实例 // 结构体声明。每个字段都是从创建作用域捕获的东西。 struct CallbackCB, F { callback: CB, map_op: F, } // 实现 ProducerCallback trait。这纯粹是样板代码。 implT, F, CB ProducerCallbackT for CallbackCB, F where F: MapOpT, CB: ProducerCallbackF::Output { type Output CB::Output; fn callbackP(self, base: P) - CB::Output where P: ProducerItemT { // 闭包的函数体在这里 let producer MapProducer { base: base, map_op: self.map_op }; self.callback.callback(producer) } } } }有点繁琐但它是有效的Map的真实实现 src/iter/map.rs#L70-L103 遵循的正是这一模式Callback结构体持有外层回调与map_op在ProducerCallback::callbackP中构造MapProducer并继续向上传递。多输入迭代器的嵌套回调以 Zip 为例当迭代器需要协调多个输入时回调会嵌套。Zip的with_producersrc/iter/zip.rs#L58-L112演示了这一点它先让a调用with_producer(CallbackA { ... })在CallbackA::callback(a_producer)里再让b调用with_producer(CallbackB { a_producer, callback })最后在CallbackB::callback(b_producer)中把两个 producer 组装成ZipProducer交给最外层的回调。每个回调结构体都只是闭包的手动版本用字段记录需要从创建作用域捕获的东西。结语Rayon 的 plumbing 层是并行迭代器的发动机舱Producer负责按需分裂和产出数据Pull 模式Consumer负责接收数据、分裂自身并用Reducer归约结果Push 模式bridge与bridge_unindexed负责在两者之间切换Splitter/LengthSplitter负责决定分裂到多细而ProducerCallback则以牺牲闭包语法糖为代价绕开了关联类型携带生命周期这一 Rust 类型系统的限制让with_producer得以实现。虽然普通用户很少直接接触 src/iter/plumbing/mod.rs其模块文档明确说明这是低层细节但对于想要编写自定义并行迭代器或组合子的开发者来说这套机制正是理解Rayon 如何把一条迭代器链变成并行任务树的关键。如果你想在此基础上继续深入可以顺着 src/iter/mod.rs 中的ParallelIterator与IndexedParallelIteratortrait 定义逐一阅读map、zip、flat_map等组合子的实现你会发现每个组合子都在这套 callback producer/consumer 的框架内各司其职。赞分享【免费下载链接】rayonRayon: A data parallelism library for Rust项目地址https://gitcode.com/gh_mirrors/ra/rayon点击查看免费下载相关推荐Rayon 线程池睡眠调度深度解析rayon-core sleep 模块的休眠协议、计数器与防死锁设计Rayon 线程池睡眠调度深度解析rayon core sleep 模块的休眠协议、计数器与防死锁设计 导读 Rayon 是一个基于工作窃取work steRay Data 内部机制深度解析执行模型、Shuffle 算法、调度与内存模型Ray Data 内部机制深度解析执行模型、Shuffle 算法、调度与内存模型 本文面向 Ray Data 的高级用户与开发者深入剖析 Ray Data人工智能分布式训练强化学习任务调度模型推理服务后端Luigi 执行模型深入解析Worker 进程内调度执行、中央调度器与外部触发机制Luigi 执行模型深入解析Worker 进程内调度执行、中央调度器与外部触发机制 Luigi 拥有一个非常简洁的执行与触发模型 执行不转移no exec任务调度工作流自动化批处理后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表