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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 07:28:08