You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

并行全局求和慢于串行?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;
                        }
                    }
                }
            }))
        }
    }
}

额外优化建议

  1. 基准测试对齐:串行版本无需创建Worker Pool,或保持与并行版本相同的线程池大小,避免无关开销干扰测试结果。
  2. 减少求和锁竞争:使用AtomicU128(需启用Rust不稳定特性)或拆分求和逻辑,替代Mutex<u128>,进一步降低线程间竞争。

内容的提问来源于stack exchange,提问作者steven-lie

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 05:25:55