Rust rayon线程池并行无性能提升问题及正确实现咨询
问题:Rust中使用Rayon线程池未获得并行性能提升
我因需要处理对速度有要求的文件开始学习Rust,写完代码后尝试实现并行执行,但使用Rayon线程池时,并行代码的运行速度慢于串行版本(示例中耗时几乎相同,实际代码中并行耗时更长)。
简化示例代码
use std::time::Instant; use std::{thread, time}; fn do_stuff(i: u64) -> u64 { let seconds = time::Duration::from_secs(i); thread::sleep(seconds); println!("{:?}", i); return i * i; } fn parallel() { let mut x = Vec::new(); let pool = rayon::ThreadPoolBuilder::new() .num_threads(8) .build() .unwrap(); for i in 1..10 { pool.install(|| { let result = do_stuff(i); x.push(result); }); } println!("{:?}", x); } fn serial() { let mut x = Vec::new(); for i in 1..10 { let result = do_stuff(i); x.push(result); } println!("{:?}", x); } fn main() { let start = Instant::now(); serial(); let duration = start.elapsed(); println!("Time elapsed (serial): {:.2}s", duration.as_secs_f64()); let start = Instant::now(); parallel(); let duration = start.elapsed(); println!("Time elapsed (parallel): {:.2}s", duration.as_secs_f64()); }
运行结果
1 2 3 4 5 6 7 8 9 [1, 4, 9, 16, 25, 36, 49, 64, 81] Time elapsed (serial): 45.00s 1 2 3 4 5 6 7 8 9 [1, 4, 9, 16, 25, 36, 49, 64, 81] Time elapsed (parallel): 45.01s
我的核心需求:
- 遍历一个包含大量元素的可迭代对象
- 将每个元素传入一个耗时计算函数
- 将函数返回的结果追加到线程启动前定义的向量中
- 使用线程池实现(不想为每个元素单独创建线程)
我有Python背景,使用ProcessPoolExecutor收集结果到列表十分简单,但在Rust中使用线程池未得到性能提升。请问我是否使用方式错误?Rust中是否有正确的实现方法?另外,我疑惑为何在parallel函数中可以对Vec执行push操作,这看起来并非线程安全的操作。
解答
1. 并行代码未提速的原因
你的parallel函数本质上还是串行执行:
rayon::ThreadPool::install()是阻塞式方法,调用后会等待传入的闭包执行完成才会继续执行下一次循环。- 循环里每次调用
pool.install()都要等当前do_stuff(i)执行完才会处理下一个i,和串行逻辑完全一致,自然不会有性能提升。
2. 正确的并行实现方式
方式一:使用Rayon的并行迭代器(推荐)
Rayon的核心优势是并行迭代器,它会自动利用线程池调度任务,无需手动管理:
use rayon::prelude::*; fn parallel() { // 直接对迭代器调用并行方法,自动分发任务到线程池 let x: Vec<u64> = (1..10).into_par_iter() .map(|i| do_stuff(i)) .collect(); println!("{:?}", x); }
这种方式会将任务并行执行,示例中总耗时会降至9秒(由耗时最长的任务决定)。
方式二:手动提交异步任务并收集结果
如果需要手动控制线程池,可以用pool.spawn()提交异步任务,通过JoinHandle收集结果:
fn parallel() { let pool = rayon::ThreadPoolBuilder::new() .num_threads(8) .build() .unwrap(); let mut handles = Vec::new(); for i in 1..10 { // spawn提交异步任务,返回用于获取结果的JoinHandle let handle = pool.spawn(move || do_stuff(i)); handles.push(handle); } // 等待所有任务完成并收集结果 let x: Vec<u64> = handles.into_iter() .map(|h| h.join().unwrap()) .collect(); println!("{:?}", x); }
3. 关于Vec::push的线程安全疑惑
你的代码能编译通过,是因为没有真正的并发写入:
- 由于
pool.install()是阻塞式的,每次只有一个线程在执行闭包,x.push(result)是串行执行的,不存在多线程同时写入Vec的情况,因此编译器没有报错。 - 如果换成真正的并发执行(比如上面的
spawn方式),直接对Vec执行push会触发编译错误——因为Vec不是线程安全的容器,多线程并发写入会导致数据竞争、内存损坏等未定义行为。
若必须在并发场景下共享写入向量,需使用线程安全容器,比如Arc<Mutex<Vec<u64>>>:
use std::sync::{Arc, Mutex}; fn parallel() { let x = Arc::new(Mutex::new(Vec::new())); let pool = rayon::ThreadPoolBuilder::new() .num_threads(8) .build() .unwrap(); for i in 1..10 { let x_clone = Arc::clone(&x); pool.spawn(move || { let result = do_stuff(i); // 加锁后再执行写入操作 let mut guard = x_clone.lock().unwrap(); guard.push(result); }); } // 等待所有任务完成(实际代码推荐用JoinHandle,这里仅做示例) thread::sleep(time::Duration::from_secs(10)); println!("{:?}", x.lock().unwrap()); }
不过这种方式会因锁竞争带来额外开销,优先推荐并行迭代器的实现方式,它会自动处理结果收集,无需手动管理线程安全。
内容的提问来源于stack exchange,提问作者dabljues
相关产品推荐
相关产品推荐

