Ch 32: std::sync::Condvar - 条件变量
std::sync::Condvar(条件变量)用于线程间的同步。一个线程可以等待某个条件为真时才继续执行,常与Mutex配合实现线程间的信号通知。
2. 基本用法
Section titled “2. 基本用法”use std::sync::{Arc, Mutex, Condvar};use std::thread;
fn main() { let pair = Arc::new((Mutex::new(false), Condvar::new())); let pair2 = Arc::clone(&pair);
// 子线程等待条件 let handle = thread::spawn(move || { let (lock, cvar) = &*pair2; let mut started = lock.lock().unwrap();
while !*started { println!("子线程等待..."); // wait会释放锁并阻塞,直到被通知 started = cvar.wait(started).unwrap(); }
println!("子线程收到通知,继续执行!"); });
// 主线程通知条件 thread::sleep(std::time::Duration::from_secs(1)); { let (lock, cvar) = &*pair; let mut started = lock.lock().unwrap(); *started = true; cvar.notify_one(); // 通知一个等待的线程 }
handle.join().unwrap();}3. 方法签名
Section titled “3. 方法签名”impl<T> Condvar { pub fn new() -> Condvar pub fn wait<T>(&self, guard: MutexGuard<T>) -> Result<MutexGuard<T>, PoisonError<MutexGuard<T>>>
pub fn wait_timeout<T>(&self, guard: MutexGuard<T>, timeout: Duration) -> Result<MutexGuard<T>, PoisonError<MutexGuard<T>>>
pub fn notify_one(&self) pub fn notify_all(&self)}4. wait - 等待条件
Section titled “4. wait - 等待条件”wait方法释放锁并阻塞,直到被唤醒:
use std::sync::{Mutex, Condvar};use std::time::Duration;
fn main() { let lock = Mutex::new(0u32); let cvar = Condvar::new();
let handle = thread::spawn(move || { let mut count = lock.lock().unwrap();
while *count < 10 { println!("等待条件... 当前: {}", *count); // wait返回时重新获取锁 count = cvar.wait(count).unwrap(); }
println!("条件满足,退出等待"); });
// 主线程增加计数并通知 for i in 1..=10 { thread::sleep(Duration::from_millis(100)); let mut count = lock.lock().unwrap(); *count = i; if i == 10 { cvar.notify_one(); } }
handle.join().unwrap();}5. wait_timeout - 超时等待
Section titled “5. wait_timeout - 超时等待”use std::sync::{Mutex, Condvar};use std::thread;use std::time::{Duration, Instant};
fn main() { let lock = Mutex::new(false); let cvar = Condvar::new();
let start = Instant::now(); let mut guard = lock.lock().unwrap();
// 等待最多2秒 let result = cvar.wait_timeout(guard, Duration::from_secs(2));
match result { Ok((new_guard, timeout_result)) => { guard = new_guard; if timeout_result.timed_out() { println!("等待超时!已过: {:?}", start.elapsed()); } else { println!("被正常唤醒,耗时: {:?}", start.elapsed()); } } Err(poisoned) => { println!("Mutex中毒: {:?}", poisoned); } }}6. notify_one - 通知一个线程
Section titled “6. notify_one - 通知一个线程”唤醒一个等待的线程:
use std::sync::{Arc, Mutex, Condvar};use std::thread;
fn main() { let queue = Arc::new((Mutex::new(Vec::new()), Condvar::new())); let queue2 = Arc::clone(&queue);
// 生产者 let producer = thread::spawn(move || { let (lock, cvar) = &*queue2; for i in 0..5 { thread::sleep(std::time::Duration::from_millis(100)); let mut q = lock.lock().unwrap(); q.push(i); println!("生产: {}", i); cvar.notify_one(); // 通知一个消费者 } });
// 消费者 let mut handles = vec![]; for _ in 0..3 { let queue3 = Arc::clone(&queue); let handle = thread::spawn(move || { let (lock, cvar) = &*queue3; loop { let mut q = lock.lock().unwrap(); while q.is_empty() { q = cvar.wait(q).unwrap(); } if let Some(item) = q.pop() { println!("消费: {}", item); if item == 4 && q.is_empty() { return; // 生产结束 } } } }); handles.push(handle); }
producer.join().unwrap(); for handle in handles { handle.join().unwrap(); }}7. notify_all - 通知所有线程
Section titled “7. notify_all - 通知所有线程”唤醒所有等待的线程:
use std::sync::{Arc, Mutex, Condvar};use std::thread;
fn main() { let flag = Arc::new((Mutex::new(false), Condvar::new())); let mut handles = vec![];
// 创建多个等待的线程 for i in 0..5 { let flag2 = Arc::clone(&flag); let handle = thread::spawn(move || { let (lock, cvar) = &*flag2; let mut started = lock.lock().unwrap();
while !*started { started = cvar.wait(started).unwrap(); }
println!("线程 {} 收到通知", i); }); handles.push(handle); }
// 主线程通知所有等待者 thread::sleep(std::time::Duration::from_millis(500)); { let (lock, cvar) = &*flag; let mut started = lock.lock().unwrap(); *started = true; println!("通知所有等待线程"); cvar.notify_all(); // 唤醒所有线程 }
for handle in handles { handle.join().unwrap(); }}8. 典型模式:生产者-消费者
Section titled “8. 典型模式:生产者-消费者”use std::sync::{Arc, Mutex, Condvar};use std::thread;use std::time::Duration;
struct ThreadPool { task_condvar: Condvar, task_mutex: Mutex<Vec<Box<dyn FnOnce() + Send>>>,}
impl ThreadPool { fn new() -> Self { ThreadPool { task_condvar: Condvar::new(), task_mutex: Mutex::new(Vec::new()), } }
fn add_task<F>(&self, task: F) where F: FnOnce() + Send + 'static, { let mut tasks = self.task_mutex.lock().unwrap(); tasks.push(Box::new(task)); self.task_condvar.notify_one(); // 通知一个工作线程 }
fn get_task(&self) { let mut tasks = self.task_mutex.lock().unwrap(); while tasks.is_empty() { tasks = self.task_condvar.wait(tasks).unwrap(); } tasks.pop() }}9. 注意事项
Section titled “9. 注意事项”- 总是使用while循环:使用while而不是if检查条件,避免虚假唤醒
- 与Mutex配合:Condvar需要与Mutex配合使用
- notify在lock之外:通常先unlock再notify,避免通知丢失
- Spurious wakeups:操作系统可能虚假唤醒,必须重新检查条件
- PoisonError:Mutex中毒时wait也会返回错误
10. 总结
Section titled “10. 总结”std::sync::Condvar核心API:
wait()- 释放锁并等待通知wait_timeout()- 带超时的等待notify_one()- 唤醒一个等待线程notify_all()- 唤醒所有等待线程
典型模式:
let (lock, cvar) = /* ... */;let mut guard = lock.lock().unwrap();while !condition { guard = cvar.wait(guard).unwrap();}// 使用guard