Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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>>。该守卫实现了 DerefDerefMut,可读写内部数据。

简单调用示例

#![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 死锁的四个必要条件

死锁必须同时满足以下四个条件:

  1. 互斥条件:资源不能被共享,只能由一个线程使用。
  2. 持有并等待条件:线程持有至少一个资源,同时等待获取其他线程持有的资源。
  3. 不可剥夺条件:资源只能由持有它的线程主动释放,不能被强制剥夺。
  4. 循环等待条件:存在一个线程循环链,每个线程都在等待链中下一个线程持有的资源。

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库提供了更轻量、功能更丰富的 MutexRwLock,生态中也有配套方式辅助做死锁检测。不过它不会自动替你消除死锁,锁顺序、作用域控制和 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 使用 LocalKeytry_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 mutthread_local!
共享性所有线程共享每个线程独立
数据竞争需要unsafe安全(因为不共享)
初始化编译期确定每个线程首次访问时初始化