rust并发编程扩展
1. 读写锁(RwLock<T>)
1.1 概述
读写锁是一种更细粒度的同步原语,它允许多个读线程同时持有锁,但只允许一个写线程独占访问。适用于读多写少的场景,可以提高并发度。
Rust的 RwLock<T>与 Mutex<T>类似,也是一个容器,包裹内部数据。通过 read()和 write()方法分别获取读锁(共享)和写锁(独占),返回的守卫离开作用域自动释放。
1.2 创建读写锁 – RwLock::new
功能:创建一个新的读写锁,内部包裹初始值 data。
接口签名:
#![allow(unused)]
fn main() {
pub fn new(data: T) -> RwLock<T>
}
- 参数:
data– 需要被保护的数据。 - 返回值:
RwLock<T>实例。
简单调用示例:
#![allow(unused)]
fn main() {
use std::sync::RwLock;
let lock = RwLock::new(5);
}
1.3 获取读锁 – RwLock::read
功能:阻塞当前线程,直到获得读锁(共享锁),返回 RwLockReadGuard。多个线程可同时持有读锁。
接口签名:
#![allow(unused)]
fn main() {
pub fn read(&self) -> LockResult<RwLockReadGuard<'_, T>>
}
- 参数:无(通过
&self获取)。 - 返回值:
LockResult<RwLockReadGuard<'_, T>>,通常调用.unwrap()获取守卫。该守卫实现了Deref,可以只读访问内部数据。
简单调用示例:
#![allow(unused)]
fn main() {
let guard = lock.read().unwrap();
println!("value = {}", *guard);
// 读锁释放
}
1.4 获取写锁 – RwLock::write
功能:阻塞当前线程,直到获得写锁(独占锁),返回 RwLockWriteGuard。若已有其他读锁或写锁,当前线程会阻塞。
接口签名:
#![allow(unused)]
fn main() {
pub fn write(&self) -> LockResult<RwLockWriteGuard<'_, T>>
}
- 参数:无。
- 返回值:
LockResult<RwLockWriteGuard<'_, T>>。该守卫实现了Deref和DerefMut,可读写内部数据。
简单调用示例:
#![allow(unused)]
fn main() {
let mut guard = lock.write().unwrap();
*guard += 1;
}
1.5 多线程共享 RwLock<T> – 配合 Arc
与 Mutex相同,RwLock也需要配合 Arc实现多线程共享所有权。
完整示例:
#![allow(unused)]
fn main() {
use std::sync::{Arc, RwLock};
use std::thread;
let data = Arc::new(RwLock::new(0));
let mut handles = vec![];
// 多个读线程
for _ in 0..3 {
let data = Arc::clone(&data);
let handle = thread::spawn(move || {
let guard = data.read().unwrap();
println!("read: {}", *guard);
});
handles.push(handle);
}
// 一个写线程
{
let data = Arc::clone(&data);
let handle = thread::spawn(move || {
let mut guard = data.write().unwrap();
*guard += 10;
println!("write: added 10");
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
}
1.6 读写锁的陷阱:写锁饥饿
在频繁读的场景下,写锁可能长时间无法获得(读锁不断被新读者获取)。Rust标准库的 RwLock实现不保证写锁优先,需要开发者注意。
2. 死锁(Deadlock)
2.1 什么是死锁
当两个或多个线程互相等待对方释放资源,导致所有线程都无法继续执行的状态,称为死锁。
2.2 死锁的四个必要条件
死锁必须同时满足以下四个条件:
- 互斥条件:资源不能被共享,只能由一个线程使用。
- 持有并等待条件:线程持有至少一个资源,同时等待获取其他线程持有的资源。
- 不可剥夺条件:资源只能由持有它的线程主动释放,不能被强制剥夺。
- 循环等待条件:存在一个线程循环链,每个线程都在等待链中下一个线程持有的资源。
2.3 Rust中常见的死锁示例
2.3.1 示例:同一线程重复获取 Mutex(递归锁问题)
Rust 的标准库 Mutex 不是递归锁,同一线程重复 lock 的行为不要依赖:标准库不保证它会成功,实际实现中可能阻塞自己,也可能 panic。
#![allow(unused)]
fn main() {
use std::sync::Mutex;
let lock = Mutex::new(0);
let _g1 = lock.lock().unwrap();
let _g2 = lock.lock().unwrap(); // 不要这样做:可能阻塞自己或 panic
}
原理:标准库
Mutex不按“同一线程可重复进入”的递归锁语义设计。需要重复进入时,应重新设计锁的作用域,或明确选择支持递归锁语义的同步原语。
2.3.2 示例:两个线程互相持有对方需要的锁
#![allow(unused)]
fn main() {
use std::sync::{Mutex, Arc};
use std::thread;
use std::time::Duration;
let a = Arc::new(Mutex::new(1));
let b = Arc::new(Mutex::new(2));
let a1 = Arc::clone(&a);
let b1 = Arc::clone(&b);
let t1 = thread::spawn(move || {
let _ga = a1.lock().unwrap();
thread::sleep(Duration::from_millis(100));
let _gb = b1.lock().unwrap(); // 等待t2释放b
});
let a2 = Arc::clone(&a);
let b2 = Arc::clone(&b);
let t2 = thread::spawn(move || {
let _gb = b2.lock().unwrap();
thread::sleep(Duration::from_millis(100));
let _ga = a2.lock().unwrap(); // 等待t1释放a
});
t1.join().unwrap();
t2.join().unwrap(); // 死锁,程序无法结束
}
2.4 Rust提供的解决死锁的方法
Rust语言层面没有自动避免死锁的机制,但标准库和生态提供了一些工具和约定来预防和检测死锁:
2.4.1 方法1:使用 try_lock避免阻塞
功能:非阻塞地尝试获取锁,如果不能立即获得则返回错误,让线程有机会做其他事或释放已有资源。
2.4.1.1 Mutex::try_lock / RwLock::try_read / RwLock::try_write
接口签名(以 Mutex::try_lock为例):
#![allow(unused)]
fn main() {
pub fn try_lock(&self) -> TryLockResult<MutexGuard<'_, T>>
}
- 参数:无。
- 返回值:
TryLockResult<MutexGuard<'_, T>>– 成功返回Ok(guard),失败返回Err(TryLockError)。
简单调用示例:
#![allow(unused)]
fn main() {
use std::sync::Mutex;
let lock = Mutex::new(0);
if let Ok(mut guard) = lock.try_lock() {
*guard += 1;
} else {
println!("锁被占用,稍后重试");
}
}
解决原理分析:通过非阻塞 try_lock,线程可以在获取失败时释放已持有的锁(通过 drop),破坏“持有并等待”条件,从而避免死锁。
2.4.2 方法2:固定锁获取顺序(避免循环等待)
通过全局约定所有线程以相同的顺序获取多个锁,可以打破循环等待条件。
示例:
#![allow(unused)]
fn main() {
// 约定:总是先锁a,再锁b
let _ga = a.lock().unwrap();
let _gb = b.lock().unwrap();
}
这样任何线程都不会出现“先锁b再锁a”的情况,循环等待被消除。
2.4.3 方法3:使用 parking_lot crate(扩展)
虽然不是标准库,但值得提及:parking_lot库提供了更轻量、功能更丰富的 Mutex和 RwLock,生态中也有配套方式辅助做死锁检测。不过它不会自动替你消除死锁,锁顺序、作用域控制和 try_lock 这类设计仍然是主要手段。
2.4.4 方法4:使用 std::sync::TryLockError模式配合超时(标准库无直接超时锁,可通过 thread::sleep配合 try_lock实现)
#![allow(unused)]
fn main() {
use std::sync::Mutex;
use std::thread;
use std::time::Duration;
let lock = Mutex::new(0);
let start = std::time::Instant::now();
loop {
if let Ok(mut guard) = lock.try_lock() {
*guard += 1;
break;
}
if start.elapsed() > Duration::from_secs(1) {
println!("超时放弃");
break;
}
thread::sleep(Duration::from_millis(10));
}
}
3. 线程构建器(std::thread::Builder)
3.1 概述
Builder允许在创建线程时配置其属性,例如线程名称、栈大小。默认使用 spawn创建线程无法设置这些属性。
3.2 创建线程构建器 – std::thread::Builder::new
功能:创建一个新的线程构建器实例,用于配置新线程的属性。
接口签名:
#![allow(unused)]
fn main() {
pub fn new() -> Builder
}
- 参数:无。
- 返回值:
Builder结构体。
简单调用示例:
#![allow(unused)]
fn main() {
use std::thread::Builder;
let builder = Builder::new();
}
3.3 设置线程名称 – Builder::name
功能:为将要创建的线程设置一个名称(主要用于调试,/proc/self/task/tid/comm下可见)。
接口签名:
#![allow(unused)]
fn main() {
pub fn name(self, name: String) -> Builder
}
- 参数:
name– 线程名称(String类型)。 - 返回值:
Builder(支持链式调用)。
简单调用示例:
#![allow(unused)]
fn main() {
let builder = Builder::new().name("my-worker-thread".to_string());
}
3.4 设置线程栈大小 – Builder::stack_size
功能:设置新线程的栈大小(字节)。默认栈大小与平台相关(通常是2MB)。
接口签名:
#![allow(unused)]
fn main() {
pub fn stack_size(self, size: usize) -> Builder
}
- 参数:
size– 栈大小(字节数)。 - 返回值:
Builder。
简单调用示例:
#![allow(unused)]
fn main() {
let builder = Builder::new().stack_size(4 * 1024 * 1024); // 4MB栈
}
3.5 创建并启动线程 – Builder::spawn
功能:使用配置好的参数创建并启动一个新线程,返回 JoinHandle。
接口签名:
#![allow(unused)]
fn main() {
pub fn spawn<F, T>(self, f: F) -> io::Result<JoinHandle<T>>
where
F: FnOnce() -> T + Send + 'static,
T: Send + 'static,
}
- 参数:
f– 线程执行闭包(与spawn相同)。 - 返回值:
io::Result<JoinHandle<T>>– 成功返回Ok(handle),失败(如栈大小非法)返回Err。
简单调用示例:
#![allow(unused)]
fn main() {
use std::thread::Builder;
let handle = Builder::new()
.name("answer-thread".to_string())
.stack_size(1024 * 1024)
.spawn(|| {
println!("Hello from named thread");
42
})
.unwrap();
let result = handle.join().unwrap();
}
3.6 完整示例
#![allow(unused)]
fn main() {
use std::thread::{Builder, current};
let builder = Builder::new()
.name("my-thread".to_string())
.stack_size(3 * 1024 * 1024);
let handle = builder.spawn(|| {
println!("Thread name: {:?}", current().name());
}).unwrap();
handle.join().unwrap();
}
4. 结构化并发(Structured Concurrency)
4.1 概念
结构化并发是一种编程范式,保证所有子线程在父作用域结束前全部完成。它不是Rust标准库的一个具体类型,而是一种编码模式,通常通过作用域线程(std::thread::scope)实现。Rust标准库从1.63版本开始支持作用域线程(scoped threads)。
4.2 作用域线程 – std::thread::scope
功能:创建一个作用域,在该作用域内生成的线程可以安全地借用作用域外的变量(无需 move)。所有作用域内线程在 scope调用返回前一定会被 join,保证没有线程泄漏。
接口签名:
#![allow(unused)]
fn main() {
pub fn scope<'env, F, T>(f: F) -> T
where
F: for<'scope> FnOnce(&'scope Scope<'scope, 'env>) -> T,
}
- 参数:
f– 一个闭包,接收一个Scope对象,在闭包内可通过Scope::spawn创建线程。 - 返回值:闭包
f的返回值。
简单调用示例:
#![allow(unused)]
fn main() {
use std::thread;
let local = vec![1, 2, 3];
thread::scope(|s| {
s.spawn(|| {
println!("first = {}", local[0]); // 可以借用local,无需move
});
s.spawn(|| {
println!("len = {}", local.len());
});
});
// 这里两个线程都已结束,local仍然有效
println!("{:?}", local);
}
4.3 作用域内创建线程 – Scope::spawn
功能:在作用域内生成一个新线程,该线程可以安全借用外部变量(生命周期受作用域限制)。
接口签名:
#![allow(unused)]
fn main() {
pub fn spawn<'scope, 'env, F, T>(&'scope self, f: F) -> ScopedJoinHandle<'scope, T>
where
F: FnOnce() -> T + Send + 'scope,
T: Send + 'scope,
}
- 参数:
f– 闭包,可以借用'scope生命周期的变量。 - 返回值:
ScopedJoinHandle<'scope, T>,可调用join()等待线程结束。
简单调用示例:
#![allow(unused)]
fn main() {
thread::scope(|s| {
let handle = s.spawn(|| {
println!("scoped thread");
100
});
let result = handle.join().unwrap();
});
}
4.4 结构化并发的优势
- 防止线程泄漏:作用域结束前自动等待所有线程。
- 允许借用外部变量:无需
move所有权,避免不必要的Arc。 - 更清晰的代码组织:父子线程关系明确。
4.5 完整示例
#![allow(unused)]
fn main() {
use std::thread;
let mut data = vec![1, 2, 3, 4];
thread::scope(|s| {
// 从data中借用切片,每个线程处理一部分
// split_at_mut 可以证明两个可变切片互不重叠
let (chunk1, chunk2) = data.split_at_mut(2);
s.spawn(move || {
for item in chunk1.iter_mut() {
*item *= 2;
}
});
s.spawn(move || {
for item in chunk2.iter_mut() {
*item *= 2;
}
});
});
// 两个线程都已结束,data被安全修改
println!("{:?}", data); // [2, 4, 6, 8]
}
5. 线程局部存储(thread_local!)
5.1 概述
线程局部存储(TLS)允许每个线程拥有变量的独立副本,互不干扰。Rust提供了 thread_local!宏来定义线程局部变量。
5.2 定义线程局部变量 – thread_local!
功能:声明一个线程局部变量,每个线程首次访问时获得一个独立初始化的实例。
宏语法结构:
#![allow(unused)]
fn main() {
thread_local! {
static NAME: Type = Expression;
// 可以有多个
}
}
static– 表示静态线程局部变量。NAME– 变量名。Type– 类型。Expression– 初始化表达式,在每个线程中独立执行。
5.3 访问线程局部变量 – with 方法
每个线程局部变量自动生成一个 with方法,用于获取该线程本地实例的引用。
功能:在当前线程上获取线程局部变量的引用,并调用传入的闭包。闭包参数是 &T。
接口签名(由宏生成,通常形式):
#![allow(unused)]
fn main() {
pub fn with<F, R>(&'static self, f: F) -> R
where
F: FnOnce(&T) -> R,
}
- 参数:
f– 接受&T并返回R的闭包。 - 返回值:闭包
f的返回值R。
简单调用示例:
#![allow(unused)]
fn main() {
use std::cell::RefCell;
thread_local! {
static COUNTER: RefCell<u32> = RefCell::new(0);
}
COUNTER.with(|c| {
*c.borrow_mut() += 1;
println!("Count: {}", *c.borrow());
});
}
5.4 使用 LocalKey的 try_with方法
功能:尝试获取线程局部变量,如果当前线程的TLS已销毁(可能在销毁期间调用),则返回错误。
接口签名(LocalKey::try_with):
#![allow(unused)]
fn main() {
pub fn try_with<F, R>(&'static self, f: F) -> Result<R, AccessError>
where
F: FnOnce(&T) -> R,
}
- 参数:
f– 闭包。 - 返回值:
Result<R, AccessError>。
简单调用示例:
#![allow(unused)]
fn main() {
COUNTER.try_with(|c| {
println!("Value: {}", *c.borrow());
}).unwrap_or_else(|_| println!("TLS already destroyed"));
}
5.5 完整示例:每个线程维护独立的计数器
#![allow(unused)]
fn main() {
use std::thread;
use std::cell::RefCell;
thread_local! {
static COUNT: RefCell<u32> = RefCell::new(0);
}
fn increment() {
COUNT.with(|c| {
*c.borrow_mut() += 1;
});
}
fn show() {
COUNT.with(|c| {
println!("Count in {:?}: {}", thread::current().id(), *c.borrow());
});
}
let t1 = thread::spawn(|| {
increment();
increment();
show(); // 输出 Count in ThreadId(1): 2
});
let t2 = thread::spawn(|| {
show(); // 输出 Count in ThreadId(2): 0
increment();
show(); // 输出 Count in ThreadId(2): 1
});
t1.join().unwrap();
t2.join().unwrap();
}
5.6 线程局部变量与普通静态变量的区别
| 特性 | static mut | thread_local! |
|---|---|---|
| 共享性 | 所有线程共享 | 每个线程独立 |
| 数据竞争 | 需要unsafe | 安全(因为不共享) |
| 初始化 | 编译期确定 | 每个线程首次访问时初始化 |