rust并发编程基础
1. 并发编程相关基础概念
1.1 并发和并行
- 并发(Concurrency):指多个任务在同一时间段内交替执行(宏观同时,微观串行)。例如单核CPU上快速切换多个线程,给人一种“同时运行”的错觉。
- 并行(Parallelism):指多个任务在同一时刻真正同时执行,需要多核CPU的支持。
💡 简单记忆:并发是“逻辑上的同时”,并行是“物理上的同时”。并发是并行的基础,但并发不一定并行。
1.2 进程和线程
- 进程:操作系统资源分配的基本单位,拥有独立的内存空间、文件句柄等。进程间通信(IPC)开销较大。
- 线程:CPU调度的基本单位,共享所属进程的内存和资源。线程间通信更高效,但也带来了数据竞争的风险。
💡 进程就像一个独立的房子,有自己的水电煤气。线程就像房子里不同房间的人,共用同一个房子的资源。
1.3 系统线程和绿色线程
- 系统线程:由操作系统内核管理,创建、切换、销毁都需要内核参与,开销相对较大。Rust标准库中的
std::thread就是系统线程。 - 绿色线程:由用户态运行时管理(如Go的goroutine),创建和切换开销极小,但需要运行时支持。Rust早期有绿色线程,但最终选择暴露系统线程作为标准。
✅ Rust线程是什么:Rust标准库直接使用1:1模型,即一个Rust线程对应一个操作系统线程。这样做的优点是:与系统交互简单、无额外运行时开销、可预测的性能。对于需要大量轻量级任务,可以选择
tokio或async-std等异步运行时(基于绿色线程思想)。
2. 线程基本操作
2.1 创建新线程 – std::thread::spawn
功能:创建一个新的OS线程,并立即开始执行传入的闭包。
接口签名:
#![allow(unused)]
fn main() {
pub fn spawn<F, T>(f: F) -> JoinHandle<T>
where
F: FnOnce() -> T + Send + 'static,
T: Send + 'static,
}
- 参数:
f– 一个闭包,在新线程中执行,闭包必须实现Send+'static(通常使用move转移所有权)。 - 返回值:
JoinHandle<T>– 代表新线程的句柄,可用于等待线程结束并获取返回值。
简单调用示例:
#![allow(unused)]
fn main() {
let handle = std::thread::spawn(|| {
println!("Hello from a new thread!");
});
}
2.2 等待线程结束 – JoinHandle::join
功能:阻塞当前线程,直到对应的线程执行完毕,并返回线程闭包的返回值。
接口签名:
#![allow(unused)]
fn main() {
pub fn join(self) -> Result<T, Box<dyn Any + Send + 'static>>
}
- 参数:无(消耗
self)。 - 返回值:
Result<T, Box<dyn Any + Send + 'static>>– 成功返回Ok(T),T是闭包返回值;失败返回Err(Box<...>)包含panic信息。
简单调用示例:
#![allow(unused)]
fn main() {
let handle = std::thread::spawn(|| 42);
let result = handle.join().unwrap(); // result = 42
}
2.3 线程闭包中的 move语义
功能:将闭包捕获的外部变量所有权转移到闭包内部,从而安全地在另一个线程中使用这些变量。
说明:不使用 move时,闭包会借用外部变量,但新线程可能存活超过变量所在作用域,导致悬垂引用。move强制转移所有权,确保变量在线程执行期间始终有效。
简单调用示例:
#![allow(unused)]
fn main() {
let data = vec![1, 2, 3];
let handle = std::thread::spawn(move || {
println!("{:?}", data); // data所有权已移入线程
});
// println!("{:?}", data); // 编译错误:data已移动
handle.join().unwrap();
}
3. 线程同步
3.1 什么是线程同步
当多个线程同时访问共享数据时,为了防止数据竞争、保证一致性和正确性,需要采用某种协调机制,这就是线程同步。例如:互斥锁、条件变量、信号量、读写锁等。
3.2 为什么需要线程同步 – 无同步的例子
考虑一个没有同步的例子:两个线程同时对一个共享计数器进行递增操作。
#![allow(unused)]
fn main() {
static mut COUNTER: i32 = 0;
let t1 = std::thread::spawn(|| {
for _ in 0..1000 {
unsafe { COUNTER += 1; }
}
});
let t2 = std::thread::spawn(|| {
for _ in 0..1000 {
unsafe { COUNTER += 1; }
}
});
t1.join().unwrap();
t2.join().unwrap();
unsafe { println!("{}", COUNTER); } // 可能不是2000,而是例如1987等随机值
}
由于线程交错执行,COUNTER += 1并非原子操作(实际上是读-改-写三步),导致丢失更新,结果不可预测。这就是典型的数据竞争,需要同步来避免。
3.3 线程同步常用手段
- 互斥锁(Mutex):保证同一时刻只有一个线程访问数据。
- 读写锁(RwLock):允许多个读线程或一个写线程同时访问。
- 条件变量(Condvar):让线程等待某个条件满足后再继续执行。
- 原子类型(Atomic):对简单类型提供无锁的原子操作。
- 消息传递(Channel):通过发送/接收消息进行线程间通信,Rust中的
mpsc是典型。
3.4 Rust互斥锁(Mutex<T>)
3.4.1 Rust的 Mutex与Java/C++的不同
- Java/C++:互斥锁通常是一个独立对象,你需要手动
lock()/unlock(),并且锁与数据是分离的。容易忘记释放锁,或者持有锁时间过长。 - Rust:
Mutex<T>是一个容器,它包裹了数据T。你无法直接访问内部数据,而必须通过lock()方法得到一个MutexGuard(类似智能指针),该守卫在作用域结束时自动释放锁。这种设计保证了数据受锁的保护,且不会忘记解锁。
✅ 核心思想:Rust的
Mutex<T>与RefCell<T>类似,都提供内部可变性——即通过不可变引用也能修改内部数据。只是RefCell在运行时检查借用规则,而Mutex通过阻塞线程来保证独占访问。
3.4.2 创建互斥锁 – Mutex::new
功能:创建一个新的互斥锁,内部包裹初始值 data。
接口签名:
#![allow(unused)]
fn main() {
pub fn new(data: T) -> Mutex<T>
}
- 参数:
data– 需要被保护的数据。 - 返回值:
Mutex<T>实例。
简单调用示例:
#![allow(unused)]
fn main() {
use std::sync::Mutex;
let m = Mutex::new(100);
}
3.4.3 加锁 – Mutex::lock
功能:阻塞当前线程,直到获得互斥锁,返回一个 MutexGuard智能指针。如果持有锁的线程panic,lock会返回错误。
接口签名:
#![allow(unused)]
fn main() {
pub fn lock(&self) -> LockResult<MutexGuard<'_, T>>
}
- 参数:无(通过
&self获取锁)。 - 返回值:
LockResult<MutexGuard<'_, T>>,通常调用.unwrap()获取MutexGuard。该守卫实现了Deref和DerefMut,可直接访问内部数据;离开作用域时自动解锁。
简单调用示例:
#![allow(unused)]
fn main() {
let guard = m.lock().unwrap(); // 获取锁
*guard += 1; // 修改内部值
println!("{}", *guard); // 读取
// guard 离开作用域,自动解锁
}
3.4.4 多线程共享 Mutex<T> – 配合 Arc
由于 Mutex<T>本身不实现 Copy,且多个线程需要共享所有权,因此通常将 Mutex<T>放入 Arc(原子引用计数)中。Arc允许多个线程同时拥有同一个 Mutex的所有权。
简单调用示例:
#![allow(unused)]
fn main() {
use std::sync::{Arc, Mutex};
use std::thread;
let counter = Arc::new(Mutex::new(0));
let mut handles = vec![];
for _ in 0..10 {
let counter = Arc::clone(&counter);
let handle = thread::spawn(move || {
let mut num = counter.lock().unwrap();
*num += 1;
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
println!("Result: {}", *counter.lock().unwrap());
}
3.5 条件变量(Condvar)
条件变量用于线程间的等待-通知机制:一个线程等待某个条件成立,另一个线程满足条件后发出通知。Rust的 Condvar总是与 Mutex配合使用。
3.5.1 创建条件变量 – Condvar::new
功能:创建一个新的条件变量。
接口签名:
#![allow(unused)]
fn main() {
pub fn new() -> Condvar
}
- 参数:无。
- 返回值:
Condvar实例。
简单调用示例:
#![allow(unused)]
fn main() {
use std::sync::{Arc, Mutex, Condvar};
let pair = Arc::new((Mutex::new(false), Condvar::new()));
}
3.5.2 条件等待 – Condvar::wait_while
功能:在持有 Mutex锁的情况下,检查等待条件(闭包)。如果闭包返回 true,表示“还需要继续等待”,它会自动释放锁并阻塞当前线程;当被其他线程通知后重新获取锁并再次检查。循环直到闭包返回 false。
接口签名:
#![allow(unused)]
fn main() {
pub fn wait_while<'a, T, F>(
&self,
guard: MutexGuard<'a, T>,
condition: F
) -> LockResult<MutexGuard<'a, T>>
where
F: FnMut(&mut T) -> bool,
}
- 参数:
guard– 已经获得的锁守卫。condition– 一个闭包,接收&mut T,返回bool。返回true表示继续等待,返回false表示停止等待并返回锁守卫。
- 返回值:返回一个新的
MutexGuard<'a, T>,此时等待条件已经不再成立。
简单调用示例:
#![allow(unused)]
fn main() {
let mut guard = lock.lock().unwrap();
guard = condvar.wait_while(guard, |data| data.is_empty()).unwrap();
// 现在 guard 中的队列非空
}
3.5.3 唤醒等待线程 – Condvar::notify_one / Condvar::notify_all
功能:唤醒一个或所有等待在该条件变量上的线程。
接口签名:
#![allow(unused)]
fn main() {
pub fn notify_one(&self)
pub fn notify_all(&self)
}
- 参数:无。
- 返回值:无。
简单调用示例:
#![allow(unused)]
fn main() {
condvar.notify_one(); // 唤醒一个线程
condvar.notify_all(); // 唤醒所有等待线程
}
3.6 多生产者单消费者通道(mpsc)
mpsc代表“Multiple Producer, Single Consumer”。它是Rust标准库提供的一个消息传递同步工具,允许多个线程发送消息,但只有一个线程接收消息。
3.6.1 创建通道 – std::sync::mpsc::channel
功能:创建一个新的异步通道,返回发送端和接收端。发送端可以克隆(多生产者),接收端独占。
接口签名:
#![allow(unused)]
fn main() {
pub fn channel<T>() -> (Sender<T>, Receiver<T>)
}
- 参数:无(通过泛型
T指定消息类型)。 - 返回值:元组
(Sender<T>, Receiver<T>)。
简单调用示例:
#![allow(unused)]
fn main() {
use std::sync::mpsc;
let (tx, rx) = mpsc::channel();
}
3.6.2 发送消息 – Sender::send
功能:将消息发送到通道。如果接收端已经关闭,返回错误。
接口签名:
#![allow(unused)]
fn main() {
pub fn send(&self, t: T) -> Result<(), SendError<T>>
}
- 参数:
t– 要发送的消息(消耗所有权)。 - 返回值:成功返回
Ok(()),失败返回Err(SendError(t))(将消息返回)。
简单调用示例:
#![allow(unused)]
fn main() {
tx.send(42).unwrap();
}
3.6.3 接收消息 – Receiver::recv
功能:阻塞当前线程,直到通道中有消息可接收,或所有发送端已关闭(此时返回错误)。
接口签名:
#![allow(unused)]
fn main() {
pub fn recv(&self) -> Result<T, RecvError>
}
- 参数:无。
- 返回值:成功返回
Ok(T),失败返回Err(RecvError)(表示没有更多消息)。
简单调用示例:
#![allow(unused)]
fn main() {
let msg = rx.recv().unwrap();
println!("Received: {}", msg);
}
3.6.4 完整示例片段
#![allow(unused)]
fn main() {
use std::sync::mpsc;
use std::thread;
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
tx.send("Hello from thread").unwrap();
});
println!("{}", rx.recv().unwrap());
}
4. 综合示例代码:生产者消费者模型实现
下面我们使用前面介绍的 Mutex、Condvar、Arc和 AtomicBool实现一个安全、可停止的生产者-消费者队列。该示例完整展示了线程同步的多项技术。
4.1 代码结构分析
Queue<T>:队列数据结构,内部包含:data: Mutex<VecDeque<T>>– 互斥锁保护的队列。condvar: Condvar– 条件变量,用于等待队列非空或停止信号。stopped: AtomicBool– 原子标志位,指示是否停止生产/消费。
push:生产者调用,添加数据并通知消费者。pop:消费者调用,如果队列为空则等待,直到有数据或停止信号。stop:设置停止标志并唤醒所有等待线程。
4.2 完整代码
#![allow(unused)]
fn main() {
use std::collections::VecDeque;
use std::sync::{Arc, Mutex, Condvar, atomic::{AtomicBool, Ordering}};
struct Queue<T> {
data: Mutex<VecDeque<T>>,
condvar: Condvar,
stopped: AtomicBool,
}
impl<T> Queue<T> {
fn new() -> Arc<Queue<T>> {
Arc::new(Queue {
data: Mutex::new(VecDeque::new()),
condvar: Condvar::new(),
stopped: AtomicBool::new(false),
})
}
fn stop(&self) {
self.stopped.store(true, Ordering::Release);
self.condvar.notify_all();
}
fn push(&self, data: T) {
let mut guard = self.data.lock().unwrap();
if self.stopped.load(Ordering::Acquire) { return; }
guard.push_back(data);
self.condvar.notify_one();
}
fn pop(&self) -> Option<T> {
let mut guard = self.data.lock().unwrap();
// 等待条件:队列非空 或者 已经停止
guard = self.condvar.wait_while(guard, |g| {
g.is_empty() && !self.stopped.load(Ordering::Acquire)
}).unwrap();
if guard.is_empty() && self.stopped.load(Ordering::Acquire) {
None
} else {
guard.pop_front()
}
}
}
pub fn demo() {
println!("...............生产者消费者模型示例开始...................");
let queue = Queue::<i32>::new();
let produce_queue = queue.clone();
let produce = std::thread::spawn(move || {
let mut count = 0;
while count < 11 {
produce_queue.push(count);
count += 1;
}
produce_queue.stop();
});
let consume_queue1 = queue.clone();
let consume1 = std::thread::spawn(move || {
while let Some(data) = consume_queue1.pop() {
println!("consumer1 consume data: {}", data);
}
});
let consume_queue2 = queue.clone();
let consume2 = std::thread::spawn(move || {
while let Some(data) = consume_queue2.pop() {
println!("consumer2 consume data: {}", data);
}
});
produce.join().unwrap();
consume1.join().unwrap();
consume2.join().unwrap();
println!("...............生产者消费者模型示例结束...................");
}
}