Skip to content

Ch 32: std::sync::Condvar - 条件变量

std::sync::Condvar(条件变量)用于线程间的同步。一个线程可以等待某个条件为真时才继续执行,常与Mutex配合实现线程间的信号通知。

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();
}
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)
}

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();
}
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);
}
}
}

唤醒一个等待的线程:

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();
}
}

唤醒所有等待的线程:

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();
}
}
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()
}
}
  1. 总是使用while循环:使用while而不是if检查条件,避免虚假唤醒
  2. 与Mutex配合:Condvar需要与Mutex配合使用
  3. notify在lock之外:通常先unlock再notify,避免通知丢失
  4. Spurious wakeups:操作系统可能虚假唤醒,必须重新检查条件
  5. PoisonError:Mutex中毒时wait也会返回错误

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