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

Rust多线程目录遍历Mutex锁阻塞及单线程工作问题求助

多线程目录遍历的协作问题

问题背景

我实现了一个遍历目录所有子目录和文件的函数,单线程版本耗时约60秒,改成4线程并行后遇到两个问题:

  • 最初用while let Some(dir) = dir_queue.lock().unwrap().pop_front()时,Mutex锁导致应用阻塞。
  • 修改为检测队列长度的循环后,锁不再阻塞,但只有1个线程在工作,其余3个线程因检测到队列空就停止了。想要实现类似队列非空 或 有线程在工作的逻辑,让多线程有效协作。

原实现代码如下:

fn get_path_files(path: &str) -> Result<Vec<String>, String> {
    let files = Arc::new(Mutex::new(Vec::new()));
    let dir_queue = Arc::new(Mutex::new(VecDeque::from([path.to_string()])));

    let mut handles = vec![];

    for _ in 0..4 {
        let files = Arc::clone(&files);
        let dir_queue = Arc::clone(&dir_queue);

        let handle = thread::spawn(move || {
            while let Some(dir) = dir_queue.lock().unwrap().pop_front() {
                if let Ok(entries) = fs::read_dir(&dir) {
                    for entry in entries {
                        if let Ok(entry) = entry {
                            println!("Dir: {:?}", entry);
                            if entry.file_type().unwrap().is_dir() {
                                dir_queue
                                    .lock()
                                    .unwrap()
                                    .push_back(entry.path().to_string_lossy().to_string());

                                print!("Dir: {:?}", entry.path().to_string_lossy().to_string());
                                files
                                    .lock()
                                    .unwrap()
                                    .push(entry.path().to_string_lossy().to_string());
                            } else {
                                files
                                    .lock()
                                    .unwrap()
                                    .push(entry.path().to_string_lossy().to_string());
                            }
                        }
                    }
                }
            }
        });

        handles.push(handle);
    }

    for handle in handles {
        handle.join().unwrap();
    }

    let result = files.lock().unwrap().clone();

    Ok(result)
}

解决方案

要解决多线程协作问题,需要跟踪两个状态:队列中是否还有待处理的目录,以及是否有线程正在处理目录。可以通过Arc<AtomicUsize>维护活跃线程计数,结合条件变量(Condvar)实现线程的等待和唤醒,避免空轮询或过早退出。

核心修改点

  • 用Arc<Condvar>配合Mutex,让线程在队列空且无活跃线程时才退出,否则等待新任务。
  • 用Arc<AtomicUsize>记录当前正在处理目录的线程数,避免队列临时为空时线程误退出。

修改后的完整代码

use std::sync::{Arc, Condvar, Mutex, atomic::{AtomicUsize, Ordering}};
use std::collections::VecDeque;
use std::thread;
use std::fs;

fn get_path_files(path: &str) -> Result<Vec<String>, String> {
    let files = Arc::new(Mutex::new(Vec::new()));
    // 队列+条件变量,用于通知线程有新任务
    let dir_queue = Arc::new((Mutex::new(VecDeque::from([path.to_string()])), Condvar::new()));
    // 活跃线程计数:记录正在处理目录的线程数量
    let active_workers = Arc::new(AtomicUsize::new(0));

    let mut handles = vec![];

    for _ in 0..4 {
        let files = Arc::clone(&files);
        let dir_queue = Arc::clone(&dir_queue);
        let active_workers = Arc::clone(&active_workers);

        let handle = thread::spawn(move || {
            loop {
                let (queue, cvar) = &*dir_queue;
                let mut queue_lock = queue.lock().unwrap();
                
                // 等待:队列空且无活跃线程时退出,否则等待新任务
                while queue_lock.is_empty() && active_workers.load(Ordering::Acquire) > 0 {
                    queue_lock = cvar.wait(queue_lock).unwrap();
                }

                // 队列空且没有活跃线程,所有任务完成,退出循环
                if queue_lock.is_empty() && active_workers.load(Ordering::Acquire) == 0 {
                    break;
                }

                // 取出一个目录,标记当前线程为活跃
                let dir = queue_lock.pop_front().unwrap();
                active_workers.fetch_add(1, Ordering::Release);
                // 解锁队列,让其他线程可以访问
                drop(queue_lock);

                // 处理目录
                if let Ok(entries) = fs::read_dir(&dir) {
                    for entry in entries {
                        if let Ok(entry) = entry {
                            let path_str = entry.path().to_string_lossy().to_string();
                            files.lock().unwrap().push(path_str.clone());

                            if entry.file_type().unwrap().is_dir() {
                                let (queue, cvar) = &*dir_queue;
                                queue.lock().unwrap().push_back(path_str);
                                // 通知等待的线程有新任务
                                cvar.notify_one();
                            }
                        }
                    }
                }

                // 处理完成,标记活跃线程数减一,通知其他线程
                active_workers.fetch_sub(1, Ordering::Release);
                dir_queue.1.notify_one();
            }
        });

        handles.push(handle);
    }

    for handle in handles {
        handle.join().unwrap();
    }

    let result = files.lock().unwrap().clone();
    Ok(result)
}

关键逻辑说明

  1. 活跃线程计数:active_workers用原子操作维护,线程开始处理目录时加1,结束时减1,确保能准确判断是否还有线程在工作。
  2. 条件变量等待/唤醒:当队列空但还有活跃线程时,线程进入等待状态;当有新目录加入队列或活跃线程数变化时,唤醒等待的线程。
  3. 锁的粒度控制:处理目录前释放队列锁,避免长时间持有锁导致其他线程阻塞;只有在修改队列或检查状态时才短暂持有锁。

内容的提问来源于stack exchange,提问作者CallMeAlien

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 07:15:08