并行全局求和慢于串行?Rust Worker Pool实现疑存逻辑问题
Worker Pool并行求和性能劣于串行的问题排查
我基于《Rust Book》实现了一个Worker Pool,使用Criterion做基准测试时,发现并行版本的全局求和远慢于串行版本,怀疑Worker Pool的实现存在逻辑问题,核心实现基本遵循《Rust Book》内容,相关代码如下:
Worker Pool实现代码
use std::sync::{Arc, Mutex}; use std::sync::mpsc::{channel, Receiver, Sender}; use std::thread::{self, JoinHandle}; pub type Job = Box<dyn FnOnce() + Send + 'static>; pub struct Woker { pub id: usize, pub thread: Option<JoinHandle<()>>, } impl Woker { pub fn new( id: usize, job_receiver: Arc<Mutex<Receiver<Job>>>, ) -> Self { Self { id, thread: Some(thread::spawn(move || { loop { let result = job_receiver.lock().unwrap().recv(); match result { Ok(job) => { job(); } Err(_) => { break; } } } })) } } } pub struct WokerPool { workers: Vec<Woker>, job_sender: Option<Sender<Job>>, } pub type ThreadSafeWorkPool = Arc<Mutex<WokerPool>>; impl WokerPool { pub fn new(size: usize) -> Self { let (job_sender, job_receiver) = channel::<Job>(); let thread_safe_job_receiver = Arc::new(Mutex::new(job_receiver)); let mut workers = Vec::with_capacity(size); for i in 0..size { workers.push(Woker::new(i,Arc::clone(&thread_safe_job_receiver))) } Self { job_sender: Some(job_sender), workers, } } pub fn execute<F>(&mut self, f: F) where F: FnOnce() + Send + 'static { self.job_sender.as_mut().unwrap().send(Box::new(f)).unwrap(); } } impl Drop for WokerPool { fn drop(&mut self) { drop(self.job_sender.take()); for worker in &mut self.workers { // println!("Work id {} stop", worker.id); if let Some(thread) = worker.thread.take() { match thread.join() { Ok(_) => {}, Err(_) => { println!("Error when join {}", worker.id) } } } } } }
基准测试代码
use criterion::{criterion_group, criterion_main, Criterion}; use learning_rust_parallel::work_pool::WokerPool; use std::sync::{Arc, Mutex}; fn criterion_benchmark(c: &mut Criterion) { let iter_time = 500000000; let size = 10; c.bench_function("parallel sum", |b| { b.iter(|| { let mut worker_pool_for_parallel = WokerPool::new(size); let sum = Arc::new(Mutex::new(0 as u128)); for _i in 0..size { let sum_clone = Arc::clone(&sum); worker_pool_for_parallel.execute(move || { let mut innner_sum = 0 ; for _index in 0..iter_time/ size { innner_sum += 1; } *sum_clone.lock().unwrap() += innner_sum }) } }) }); c.bench_function("sequential sum", |b| { b.iter(|| { // create same size pool to make sequential version contain time of create threads. let _worker_pool_for_parallel = WokerPool::new(5); let mut sum: u128 = 0; for _i in 0..iter_time { sum += 1; } }) }); } criterion_group!(benches, criterion_benchmark); criterion_main!(benches);
问题分析与修复
核心问题:Mutex锁定范围过大
Worker线程的核心逻辑中,job_receiver.lock().unwrap().recv()会在调用recv()的整个阻塞期间持有Mutex锁。这导致所有Worker线程必须排队等待同一个锁,任务被串行执行,完全无法发挥并行优势。
修复方案
调整锁的持有范围,仅在接收Job的瞬间持有锁,拿到Job后立即释放锁,让其他Worker可以获取任务:
impl Woker { pub fn new( id: usize, job_receiver: Arc<Mutex<Receiver<Job>>>, ) -> Self { Self { id, thread: Some(thread::spawn(move || { loop { // 用代码块限制锁的持有范围 let job = { let receiver = job_receiver.lock().unwrap(); receiver.recv() }; match job { Ok(job) => { job(); } Err(_) => { break; } } } })) } } }
额外优化建议
- 基准测试对齐:串行版本无需创建Worker Pool,或保持与并行版本相同的线程池大小,避免无关开销干扰测试结果。
- 减少求和锁竞争:使用
AtomicU128(需启用Rust不稳定特性)或拆分求和逻辑,替代Mutex<u128>,进一步降低线程间竞争。
内容的提问来源于stack exchange,提问作者steven-lie
相关产品推荐
相关产品推荐

