☰
Rayon 1.12 数据并行编程实战指南:从 par_iter 到自定义线程池与 WebAssembly 支持
2026/10/7 7:59:20 网站建设 项目流程

【免费下载链接】rayon

Rayon: A data parallelism library for Rust

项目地址:https://gitcode.com/gh_mirrors/ra/rayon
点击查看免费下载

导读

本文基于 Rayon 项目官方 README 与仓库源码,系统讲解 Rust 数据并行库 Rayon(当前仓库版本 1.12.0)的核心用法:如何用par_iter()一行代码把顺序迭代改成并行迭代、如何用join与scope手动拆分任务、如何通过ThreadPoolBuilder定制线程池,以及它在无线程 WebAssembly 目标上的降级与多线程适配方案。读完本文,你将掌握将既有顺序计算转换为并行计算的标准套路,并能理解其"数据竞争自由"保证的边界条件。

Rayon 是什么

Rayon 是 Rust 生态中一个数据并行(data-parallelism)库。它的定位非常明确:极轻量,让开发者把一段顺序计算改造成并行计算的成本降到最低——通常只需把foo.iter()改成foo.par_iter()。更重要的是,它保证数据竞争自由(data-race freedom),这意味着大多数并行 bug 在编译期就被排除掉了:只要代码能编译通过,程序的行为通常与改写前保持一致。

从仓库源码看,项目由两个 crate 组成(见 Cargo.toml):

  • rayon(1.12.0):面向用户的公共 API 层,提供并行迭代器、切片排序、集合扩展等高级构造;
  • rayon-core(1.13.0):核心运行时,负责线程池注册表(registry)、任务窃取、作用域与惰性同步(latch)等底层实现,其稳定 API 已镜像到rayon中。

快速上手:把顺序迭代器变成并行迭代器

Rayon 最有吸引力的一点是改造成本极低。官方 README 给出的示例非常经典——计算平方和:

use rayon::prelude::*; fn sum_of_squares(input: &[i32]) -> i32 { input.par_iter() // <-- 只需改这一行! .map(|&i| i * i) .sum() }

改动仅在一处:iter()变成par_iter(),其余链式调用(map、sum)原样保留,返回值与顺序版本完全一致。

并行迭代器如何工作

从源码结构看,并行迭代器的核心机制位于 src/iter/ 模块。它围绕两个关键 trait 组织(见 src/lib.rs 的说明):

  • ParallelIterator:定义所有并行迭代器的通用方法,如map、for_each、filter、fold等(对应 src/iter/mod.rs);
  • IndexedParallelIterator:为支持随机访问的迭代器(如切片、向量)追加能力,是zip、enumerate、按索引切分等高阶操作的基础。

并行迭代器的职责是决定如何把数据划分成任务,并在运行时动态自适应以获得最大性能:它不会机械地按固定块大小切分,而是根据当前的工作负载与可用线程数动态决定拆分粒度。这正是它与"手动把数据均分到 N 个线程"这类朴素方案的本质区别。

prelude:一条 import 引入全部并行 trait

并行迭代器涉及多个 trait,手工逐个导入既繁琐又易错。Rayon 的解决方案是把它们统一打包进 prelude(见 src/prelude.rs):

use rayon::prelude::*;

这一行导入包括ParallelIterator、IndexedParallelIterator、IntoParallelIterator、IntoParallelRefIterator、IntoParallelRefMutIterator、FromParallelIterator、ParallelExtend、ParallelSlice、ParallelSliceMut、ParallelString等在内的一整套 trait(src/prelude.rs)。建议在每一个用到并行迭代器 API 的模块顶部都加上这一行。

从库文档看更多高级构造

除了par_iter,src/lib.rs 还列出了其他两类开箱即用的高级能力:

  • 并行排序:par_sort方法(ParallelSliceMuttrait 提供)可以对&mut [T]切片或向量做并行排序,其实现见 src/slice/sort.rs,并有对应的 panic 安全测试 tests/sort-panic-safe.rs;
  • 并行扩展集合:par_extend(ParallelExtendtrait 提供)可以高效地把并行迭代器产出的元素批量写入集合,测试见 tests/collect.rs。

无数据竞争:编译通过即行为一致

并行编程最令人头疼的是各种诡异的数据竞争 bug。Rayon 的 API全部保证数据竞争自由,这排除了绝大多数并行 bug(虽然并非全部)。换句话说:只要你的代码能编译通过,它通常做的还是原来那件事。

不过 README 明确指出了两个需要注意的边界:

  1. 副作用顺序不可保证:如果迭代器带有副作用(例如通过 Rust channel 向其他线程发消息、或写磁盘),这些副作用发生的顺序可能与顺序版本不同。仓库测试 tests/iter_panic.rs 等用例也印证了并行执行顺序的动态性。
  2. 存在更高性能的替代方法:某些场景下,并行迭代器会提供顺序迭代方法的替代版本,性能更高(例如reduce之于fold、try_reduce之于try_fold),取舍时需要按实际语义选择。

引入 Rayon:Cargo.toml 与版本要求

在 Cargo.toml 中加入依赖即可:

[dependencies] rayon = "1.12"

仓库当前版本为 1.12.0(见 Cargo.toml),工作区还包含rayon-demo与rayon-core两个成员(Cargo.toml)。

版本要求:Rayon 目前要求rustc 1.85.0或更高版本。仓库的 workspace 配置印证了这一点——rust-version = "1.85"、edition = "2024"(Cargo.toml)。使用旧版工具链将无法编译。

另外要注意:由于rayon-core需要保证全局线程池的唯一性与线程池间的协调,它禁止在同一目标里链接多个版本的自己。如果不同 crate 对rayon-core的版本约束(如过紧的~或不等式约束)相互冲突,你会看到类似下面的构建错误,需要先解决版本冲突:

error: native library `rayon-core` is being linked to by more than one package, and can only be linked to by one package

超越迭代器:join 与 scope

并行迭代器虽然好用,但并非万能。当你需要更多灵活性时,Rayon 提供了两个底层原语(源码位于 rayon-core/src/join/mod.rs 与 rayon-core/src/scope/mod.rs):

  • join(a, b):接收两个闭包并"可能并行"地执行它们,返回二元组结果。它类似"开两个线程各跑一个闭包",但实现完全不同、开销极低;
  • scope(|s| ...):创建一个作用域,在其中可以任意创建并行任务(含嵌套任务与嵌套 scope),scope 会一直存活到其中所有任务完成。

join 背后的工作窃取原理

从 rayon-core/src/join/mod.rs 的实现可以看到 join 的完整执行路径:

  1. 当前线程把任务 B 包装为StackJob并压入自己的本地工作队列(deque),同时用SpinLatch记录其完成状态;
  2. 当前线程先执行任务 A——此时任务 B 处于"可被偷取"的公告状态;
  3. 任务 A 完成后,当前线程尝试从本地栈弹出任务 B 直接执行;若任务 B 已被其他线程窃取(steal),则当前线程转而寻找其他工作,并在任务 B 完成前阻塞等待(wait_untillatch)。

这套机制就是著名的**工作窃取(work stealing)**策略:Rayon 运行时维护一个固定大小的 worker 线程池,只有在存在空闲 CPU 时才会真正并行执行代码;否则join的两个闭包会在同一个线程上顺序执行。这保证了即使线程池繁忙,程序也不会过度并行化。

调用join时有两点注意事项(源码文档中有明确警告):

  • 闭包应近似 CPU 密集:若闭包内含可能阻塞的 I/O(如等待网络请求),整体性能可能很差;若一个闭包阻塞等待另一个闭包(例如通过 channel),甚至可能死锁;
  • panic 传播规则:无论如何两个闭包都会被执行。若单个闭包 panic,join以相同 panic 值传播;若两个都 panic,则以第一个闭包的 panic 值传播。若任务 A panic,实现还会先等待任务 B 完成再恢复展开,因为任务 B 可能持有外层栈帧的引用(见 join/mod.rs 的join_recover_from_panic)。

从 join 到完整任务树:quick_sort 示例

join最典型的应用是递归分治。以下示例(来自 rayon-core/src/join/mod.rs 的文档)用 join 实现快速排序——注意这并非最优实现,实际排序请优先使用par_sort:

use rayon::prelude::*; let mut v = vec![5, 1, 8, 22, 0, 44]; quick_sort(&mut v); assert_eq!(v, vec![0, 1, 5, 8, 22, 44]); fn quick_sort<T: PartialOrd + Send>(v: &mut [T]) { if v.len() > 1 { let mid = partition(v); let (lo, hi) = v.split_at_mut(mid); rayon::join(|| quick_sort(lo), || quick_sort(hi)); } } // partition 把所有 <= 枢轴(取切片最后一个元素)的项 // 交换到前半部分,并返回枢轴所在的分界点下标 fn partition<T: PartialOrd + Send>(v: &mut [T]) -> usize { let pivot = v.len() - 1; let mut i = 0; for j in 0..pivot { if v[j] <= v[pivot] { v.swap(i, j); i += 1; } } v.swap(i, pivot); i }

每次递归调用把当前区间一分为二交给join,两个子区间可以在不同 worker 线程上并行排序,整个调用树便形成一棵并行任务树——这正是工作窃取调度器最擅长的负载形态。仓库的编译失败测试(如 rayon-core/src/compile_fail/quicksort_race1.rs)也从反面验证了这类分治写法对线程安全边界的严格要求。

更精细的控制:自定义线程池与全局池配置

默认情况下 Rayon 使用一个全局线程池,大多数场景无需关心它的存在。但当你需要控制线程数量、线程名、栈大小等参数时,可以用ThreadPoolBuilder(定义见 rayon-core/src/lib.rs)。

创建独立线程池

let pool = rayon::ThreadPoolBuilder::new().num_threads(22).build().unwrap();

配置全局线程池

rayon::ThreadPoolBuilder::new().num_threads(22).build_global().unwrap();

build_global的初始化是可选的——不调用的话,全局池会在首次使用时以默认配置自动初始化。全局池恰好初始化一次,之后配置不可更改,再次调用build_global会返回错误。官方建议仅在两种场景下手动调用:

  1. 你想改变默认配置;
  2. 你在跑基准测试——预先初始化可以让 worker 线程就绪,使第一轮迭代的数据更稳定(但收益很小)。

ThreadPoolBuilder 全部配置项

从 rayon-core/src/lib.rs 的默认实现与各 setter 方法可以看出完整配置面:

方法作用默认行为
num_threads(n)设置线程数;传 0 或不调用则由运行时自动选择优先读RAYON_NUM_THREADS环境变量,否则取逻辑 CPU 数(std::thread::available_parallelism,失败则回退为 1)
thread_name(f)按线程下标(usize)生成线程名无自定义名
stack_size(n)设置 worker 线程栈大小交给std::thread::Builder默认值
panic_handler(h)兜底处理无法向上传播的 panic(主要针对spawnAPI)未设置时直接中止进程,原则是"panic 不应被无视";handler 自身 panic 也会中止进程,建议内部包一层catch_unwind
start_handler(h)/exit_handler(h)线程启动/退出回调,参数为线程下标,可能在多线程中并行调用无
spawn_handler(f)自定义线程创建逻辑(返回io::Result),可结合std::thread::scope做作用域线程默认设置名称与栈大小并传播错误
use_current_thread()把调用线程并入池,保证其下标为 0;该线程不跑主工作窃取循环关闭
breadth_first()(已弃用)提示 worker 从本地队列底部取任务做广度优先执行默认深度优先;建议改用scope_fifo/spawn_fifo

线程数解析顺序的源码细节

rayon-core/src/lib.rs 中get_num_threads的解析顺序值得注意:

  1. num_threads(n)指定了非 0 值 → 直接采用;
  2. 否则读RAYON_NUM_THREADS环境变量:正整数直接采用,为 0 则回退到默认;
  3. 再回退到已弃用的RAYON_RS_NUM_CPUS(新变量是它的一对一替代,两者同时存在时优先RAYON_NUM_THREADS);
  4. 都不满足时用std::thread::available_parallelism()(失败则 1)。

此外,单个线程池的线程数存在上限max_num_threads(),由睡眠计数器的AtomicUsize位宽限制决定(见 rayon-core/src/lib.rs),超出部分会被削减。current_num_threads()则返回当前注册表的线程数(可用于指导拆分次数,并行迭代器内部就使用它)。

scope 的变体与惰性同步

如果只是需要在当前线程内创建任务(而不是先install到其他池),可用scope/in_place_scope;需要 FIFO 调度语义时用scope_fifo/in_place_scope_fifo(均从 rayon-core/src/lib.rs 再导出到rayon)。仓库的 scope 测试见 rayon-core/src/scope/test.rs,spawn 相关测试见 rayon-core/src/spawn/test.rs。

WebAssembly 支持:无线程目标的降级与多线程适配

Rayon 依赖std的线程 API,但部分目标平台的线程实现不完整、总是返回Unsupported错误,wasm32-unknown-unknown与wasm32-wasi就是典型例子。针对这些目标,Rayon 配置了全局回退模式(global fallback),而不是在创建隐式全局池时 panic(详见 rayon-core/src/lib.rs 的文档)。

默认行为:顺序回退

构建到 WebAssembly 时,Rayon默认视其为无多线程平台,回退到顺序迭代。你的既有代码无需任何修改即可编译并运行成功——代价是只用单核、运行较慢。

回退模式的具体语义(类似把RAYON_NUM_THREADS设为 1,但又不完全相同):

  • join的两个闭包会顺序执行(没有其他线程可分工作);
  • 由于池不独立于主线程运行,非阻塞调用如spawn可能根本不会执行,除非broadcast这类低优先级调用给了它们执行窗口;
  • 回退模式不模拟线程抢占或async任务切换,但yield_now/yield_local可以主动让出执行时间;
  • 注意:显式调用ThreadPoolBuilder的方法不做任何回退,会直接报告错误。

开启多线程:需要适配器与项目配置

若要在 Web 上获得真正的多线程并行,需要借助适配器并调整项目配置,以弥合 WebAssembly 线程与常规平台线程的差异。官方推荐的是wasm-bindgen-rayon(其完整配置方案见该适配器文档,此处不做展开)。同时,仓库 Cargo.toml 提供了一项与 Web 相关的特性:

[features] # 此特性在浏览器的 main 线程上切换到自旋锁实现, # 以避开被禁止的 `atomics.wait`。 # 仅对 `wasm32-unknown-unknown` 目标有用。 web_spin_lock = ["dep:wasm_sync", "rayon-core/web_spin_lock"]

web_spin_lock启用后,rayon-core会改用wasm_sync提供的同步原语(见 rayon-core/src/lib.rs 的条件编译),使得在浏览器主线程上也能安全等待,而不触碰被 WebAssembly 规范禁止的atomics.wait。

亲自体验:rayon-demo 演示程序

仓库自带一个演示/基准集合rayon-demo,是感受 Rayon 并行效果最直接的方式。README 给出的 N 体模拟可视化命令如下:

> cd rayon-demo > cargo run --release -- nbody visualize

可视化窗口中按s切换到顺序执行、按p切换到并行执行,可以直观对比两者的运行速度差异。

rayon-demo的命令行入口定义在 rayon-demo/src/main.rs,支持以下演示与基准(每个子命令可用--help查看自身参数):

名称内容
life康威生命游戏(Conway's Game of Life)
nbody多体物理模拟(bench跑基准并输出耗时,visualize可视化)
sieve埃拉托斯特尼筛法求素数
matmul并行矩阵乘法
mergesort并行归并排序
noop启动空任务以测量 CPU 开销
quicksort并行快速排序
tsp旅行商问题求解器(示例数据集在 rayon-demo/data/tsp)

例如 nbody 基准(源码见 rayon-demo/src/nbody/mod.rs):

> cd rayon-demo > cargo run --release -- nbody bench --mode par --bodies 4000 --ticks 100

--mode可取par(并行)、parreduce(并行 + reduce 归约)、seq(顺序),默认三种都跑,并输出并行耗时、顺序耗时与加速比(speedup):

Parallel time : ... ns ParReduce time : ... ns Sequential time : ... ns Parallel speedup : ...

完整基准套件还可通过cargo bench或rayon-demo bench运行;更多用法运行cargo run --release -- --help查看。

常见问题与其他资源

如果使用过程中遇到疑问,仓库根目录的 FAQ.md 汇集了常见问题解答,涵盖并行迭代器语义、线程池行为、与标准库迭代器的差异等话题。

仓库中还包含丰富的测试资产可作深入参考:

  • 集成测试目录 tests/:涵盖跨线程池调用(cross-pool.rs)、并行收集(collect.rs)、字符串处理(str.rs)、命名线程(named-threads.rs)等场景;
  • 编译失败测试 src/compile_fail/:用trybuild风格的负向测试验证"非 Send/非 Sync 数据无法进入并行迭代器"等编译期约束(如 cell_par_iter.rs、rc_par_iter.rs)。

许可证

Rayon 以 MIT 与 Apache License 2.0 双许可证发布,详见仓库根目录的 LICENSE-APACHE 与 LICENSE-MIT。向项目提交 PR 即视为同意这些许可条款。

小结

Rayon 的价值在于把"并行化"从一件需要精心设计的事,变成一行代码的改动:iter()→par_iter()。它的数据竞争自由保证让并行代码的安全边界由编译器把关,而join/scope与ThreadPoolBuilder则为需要精细控制的任务树与线程池形态提供了完整的进阶手段。若你的目标平台是 WebAssembly,默认的顺序回退保证了代码可移植,而wasm-bindgen-rayon适配器 +web_spin_lock特性则为多线程 Web 场景保留了升级路径。

【免费下载链接】rayon

Rayon: A data parallelism library for Rust

项目地址:https://gitcode.com/gh_mirrors/ra/rayon
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询