📚 Rust 课程系列

  1. 课程概览
  2. 基础语法(一):变量、数据类型与字符串
  3. 切片(Slice):序列的借用视图
  4. 基础语法(二):运算符、表达式与控制流
  5. 函数与输入输出
  6. 所有权、借用与生命周期
  7. 结构体
  8. 枚举
  9. 模式匹配
  10. 类型系统:泛型、trait 与多态
  11. 集合与容器
  12. 错误处理与 Panic 恢复
  13. 模块、属性与宏
  14. 智能指针、迭代器与闭包
  15. 并发与异步编程(本文)
  16. Unsafe Rust 与常用 trait 详解
  17. 工具链、Cargo 与外部 crate
  18. 最佳实践、性能与调试

Rust 的并发安全由类型系统在编译期保证:Send/Sync 标记 trait 让「能否跨线程传递/共享」成为类型的一部分,从而「无畏并发」(fearless concurrency)。本章按主题组织:并发安全基础(Send/Sync)、线程(std::thread、消息通道、scope 线程、线程本地存储)、共享状态与同步原语(Mutex/RwLock/CondvarBarrierOnceLock),以及基于 async/.await 的异步编程(Future、运行时、join!/select!)。

并发安全:Send 与 Sync

trait含义自动推导规则反例
SendT 的所有权可安全跨线程转移所有字段均为 SendRc<T>*const T
Sync&T 可安全跨线程共享所有字段均为 Sync(等价于 &T: SendCell<T>RefCell<T>

💡 提示Send/Sync 是标记 trait(marker trait),没有方法,仅作编译器的「通行证」。手动 impl Send/impl Sync 是 unsafe 的——你在向编译器担保安全性。

🔄 对比:C++ 线程安全靠程序员自觉(std::mutex 配对靠约定),Rust 在编译期拒绝不安全的跨线程访问——Rc 不是 Send,编译器直接报错,而非运行时崩溃。

线程

线程创建与消息传递

1
2
3
4
5
6
7
8
9
use std::thread;
use std::sync::mpsc;

let (tx, rx) = mpsc::channel();
let handle = thread::spawn(move || {
tx.send(42).unwrap(); // move tx into thread
});
println!("Received: {}", rx.recv().unwrap());
handle.join().unwrap();

💡 提示thread::spawn 要求闭包 'static,因此捕获的局部变量必须 move 进线程。若需借用栈上数据,见下方 Scope 线程

💡 提示mpsc 适合生产者-消费者模式;性能敏感场景可用 crossbeam 提供更高性能的多生产者通道。

💡 提示join 返回 Result<T, Box<dyn Any + Send>>–子线程 panic 被捕获为 Err,不会传染调用线程。这是线程 panic 隔离的体现;panic = "abort" 模式下则直接终止进程。

move 闭包与所有权转移

thread::spawn 要求闭包满足 F: Send + 'static。闭包默认按引用捕获变量,而局部变量引用的生命周期受限于当前栈帧——线程可能比调用者活得更久,因此编译器拒绝不带 move 的借用:

1
2
3
4
5
6
let v = vec![1, 2, 3];
// thread::spawn(|| println!("{:?}", v));
// error[E0373]: closure may outlive the current
// function, but it borrows `v`
let handle = thread::spawn(move || println!("{:?}", v));
// v 的所有权已转移,此处不能再使用 v

要点:

  • move 把捕获变量的所有权转移进闭包,闭包随之成为 'static,可以安全送往新线程。
  • 转移后原变量失效;Copy 类型(如 i32)不受影响,按位复制即可。
  • 多线程共享同一数据时,惯用法是先 Arc::clonemove——克隆的是句柄(引用计数),数据本身仍只有一份:
1
2
3
4
5
6
7
8
9
10
11
use std::sync::{Arc, Mutex};

let counter = Arc::new(Mutex::new(0));
let mut handles = vec![];
for _ in 0..4 {
let c = Arc::clone(&counter); // 每个线程克隆一个句柄
handles.push(thread::spawn(move || {
*c.lock().unwrap() += 1; // move 转移句柄 c
}));
}
for h in handles { h.join().unwrap(); }

⚠️ 注意move 只决定闭包的捕获方式,不影响闭包体内的借用规则。若仅需借用栈上数据且能保证线程先于栈帧结束,使用 Scope 线程 可免去 move

🔄 对比:C++ 线程靠 lambda 捕获列表区分 [=]/[&],误用 [&] 捕获局部变量是悬垂引用的经典来源;Rust 用 'static 约束在编译期强制你显式 move 或改用 scope 线程。

线程配置:Builder

1
2
3
4
5
6
7
8
use std::thread;

let handle = thread::Builder::new()
.name("worker".into())
.stack_size(8 * 1024 * 1024)
.spawn(|| { /* ... */ })
.unwrap();
handle.join().unwrap();

💡 提示:命名线程配合日志与 panic 信息可快速定位问题;thread::current().name() 返回 Option<&str>thread::current().id() 返回 ThreadId 用于区分工作线程。

⚠️ 注意thread::spawn 创建失败时 panic;Builder::spawn 返回 io::Result<JoinHandle>,可优雅处理资源耗尽等错误。默认栈约 2 MiB;工作线程数可参考 thread::available_parallelism() 返回的可用并行度。

线程让出与阻塞

1
2
3
4
use std::{thread, time::Duration};

thread::sleep(Duration::from_millis(100));
thread::yield_now();

💡 提示yield_now 仅为「提示」–调度器可能忽略并继续运行当前线程。忙等场景应优先 sleep 或阻塞原语让出 CPU,而非空转。

park 与 unpark

1
2
3
4
5
6
7
8
9
10
11
use std::{thread, time::Duration};

let handle = thread::spawn(|| {
println!("parking...");
thread::park();
println!("resumed");
});

thread::sleep(Duration::from_millis(50));
handle.thread().unpark();
handle.join().unwrap();

💡 提示unpark 采用「许可」语义–若目标线程尚未 park,许可被保存,下次 park 立即返回,因而不会丢失唤醒信号;Condvar::notify 在无等待者时则信号丢失。

🔄 对比park/unpark 适合「一对一」唤醒(已知目标 Thread 句柄);涉及条件判断与多线程协调仍应使用 Condvar+Mutex。Java 的 LockSupport.park/unpark 同为许可模型。

Scope 线程

1
2
3
4
5
6
7
8
use std::thread;
let mut numbers = vec![1, 2, 3, 4];
thread::scope(|s| {
s.spawn(|| println!("{:?}", numbers)); // borrow &numbers
s.spawn(|| numbers.push(5)); // borrow &mut numbers
});
// scope exits: all threads joined, numbers usable again
numbers.push(6);

💡 提示thread::scope 借助作用域保证所有线程在作用域结束时必然 join,因此允许借用栈上局部数据,避免把数据强行 Arc/move 的繁琐。

⚠️ 注意:scope 内借用规则仍然生效——编译器会检查同一作用域内的并发借用是否冲突。上例中两个 spawn 的借用不重叠(一个 &T,一个 &mut T),编译器会确保它们不会同时活跃。

线程本地存储(Thread Local Storage)

1
2
3
4
5
use std::cell::RefCell;
thread_local! {
static COUNTER: RefCell<u32> = RefCell::new(0);
}
COUNTER.with(|c| *c.borrow_mut() += 1); // each thread has its own copy

💡 提示:TLS 让每个线程拥有独立副本,无需同步即可修改,常用于缓存、随机数生成器、请求上下文等「每线程一份」的状态。因 thread_local! 中的值不是 Sync,内部用 RefCell 而非 Mutex

共享状态与同步原语

Mutex 与 Arc

1
2
3
4
5
6
7
8
use std::sync::{Arc, Mutex};
let counter = Arc::new(Mutex::new(0));
let c2 = Arc::clone(&counter);
let handle = std::thread::spawn(move || {
*c2.lock().unwrap() += 1; // lock() returns MutexGuard (RAII)
});
handle.join().unwrap();
println!("{}", *counter.lock().unwrap()); // 1

⚠️ 注意:若只是计数,优先使用 AtomicUsize 等原子类型,避免锁开销。Mutex 适用于保护复杂共享数据。

⚠️ 注意:避免在持锁期间调用外部代码或再加锁——这是死锁的经典来源,尽量缩小锁的持有范围。

RwLock

1
2
3
4
5
6
7
use std::sync::RwLock;
let lock = RwLock::new(5);
{
let r1 = lock.read().unwrap();
let r2 = lock.read().unwrap();
} // 多读者并发
{ *lock.write().unwrap() += 1; } // 独占写

⚠️ 注意std::sync::RwLock 在写饥饿场景下可能表现不佳;性能敏感场景考虑 parking_lot::RwLock(写优先,内存开销更小)。

Condvar

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
use std::sync::{Arc, Mutex, Condvar};
let pair = Arc::new((Mutex::new(false), Condvar::new()));
let pair2 = Arc::clone(&pair);

std::thread::spawn(move || {
let (lock, cvar) = &*pair2;
let mut started = lock.lock().unwrap();
*started = true;
cvar.notify_one();
});

let (lock, cvar) = &*pair;
let mut started = lock.lock().unwrap();
while !*started {
started = cvar.wait(started).unwrap(); // 释放锁并等待唤醒
}

💡 提示Condvar::wait 释放 MutexGuard 并挂起线程,被唤醒后重新获取锁。必须用 while 循环检查条件(防止虚假唤醒)。

🔄 对比:Condvar 用法与 C++ 的 std::condition_variable、Java 的 Object.wait()/notify() 几乎一致。

OnceLock 与 OnceCell(Rust 1.70+)

1
2
3
4
5
6
7
use std::sync::OnceLock;

static CONFIG: OnceLock<Config> = OnceLock::new();

fn get_config() -> &'static Config {
CONFIG.get_or_init(|| Config { /* ... */ })
}
类型线程安全用途
OnceLockstatic 声明的全局延迟初始化
OnceCell单线程场景的延迟初始化

关键 API:get_or_init(首次访问时初始化)、set(尝试设置,返回 Result)、get_or_try_init(支持返回 Result 的初始化)。

💡 提示OnceLock 取代了过去的 lazy_static/once_cell 第三方方案,是延迟初始化全局变量的标准做法。

Barrier

1
2
3
use std::sync::{Arc, Barrier};
let barrier = Arc::new(Barrier::new(4)); // 4 threads must all arrive
// each thread: barrier.wait(); // blocks until all 4 reach here

🔄 对比:Barrier ≈ C++ 的 std::latch/std::barrier(C++20)、Java 的 CyclicBarrier

异步编程

async/.await 基础

1
2
3
4
5
#[tokio::main]
async fn main() {
let val = async { 42 }.await;
println!("{}", val);
}
  • async fn 返回实现 Future 的匿名类型;.await 挂起当前任务,让出线程。
  • 并发原语:join!(并发等待多个 Future)、select!(多路复用)、spawn(任务调度)。

💡 提示:Rust 的异步是「协作式」的——Future 只有被驱动(poll)才会推进,.await 点是让出执行权的地方。运行时(如 tokio)负责调度任务、在 IO 就绪时唤醒它们。async fn 本身不执行任何代码,只有 .await 或交给运行时 spawn 后才运行。

🔄 对比:Go 的 goroutine 是抢占式调度(运行时自动切换),Rust 的 async 是协作式。好处是切换开销极低且可预测,代价是须避免在 .await 之间执行长时间阻塞操作。

⚠️ 注意:异步块默认单线程调度,CPU 密集任务应使用 spawn_blocking 或专用线程池,否则会饿死其他任务。

运行时与生态

  • 应用层:使用 tokio 或 async-std 作为运行时。
  • 库层:尽量不直接依赖特定运行时以提高可组合性(使用 futures 兼容 trait)。
  • IO 密集型服务:HTTP、数据库访问、网络爬虫等使用 async 可获得高并发吞吐。

⚠️ 注意:在 async 上下文中调用 std::thread::sleepstd::net::TcpListener::accept 等同步阻塞操作会阻塞整个执行线程,导致同线程其他任务全部卡住。务必使用运行时提供的异步版本(如 tokio::time::sleep)。

Future trait(进阶)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};

struct MyFuture { count: u32 }

impl Future for MyFuture {
type Output = u32;
fn poll(
mut self: Pin<&mut Self>,
_cx: &mut Context<'_>,
) -> Poll<Self::Output> {
if self.count > 0 { self.count -= 1; Poll::Pending }
else { Poll::Ready(42) }
}
}

🔬 进阶Future::poll 是异步的底层机制。返回 Poll::Pending 时,必须通过 cx.waker() 安排唤醒,否则任务将永远不会再被 poll。上例省略了 waker 注册,仅作演示——实际手写 Future 必须正确处理 waker。

🔬 进阶Pin 保证自引用 Future 不会被移动,这是 Rust 异步零成本抽象的关键。编译器生成的 Future 状态机可能包含自引用,Pin 在类型层面阻止 mem::swap/mem::replace 等移动操作。

异步 trait

自 Rust 1.75 起,trait 中直接写 async fn 已稳定:

1
2
3
trait AsyncService {
async fn process(&self, data: &str) -> String;
}

⚠️ 注意:若需把异步 trait 当作 dyn Trait 动态分派,原生 async fn 在 trait 中返回不透明 Future 类型,无法直接 Box<dyn Future>。此时仍需借助 async-trait crate(把 async fn 转为返回 Pin<Box<dyn Future>>)。

🔬 进阶async fn in trait 的动态分派限制源于 Rust 类型系统——每个 async fn 返回的匿名 Future 类型大小不同,无法统一放入 trait object 虚表。async-trait 通过 Box 擦除类型绕过,代价是一次堆分配,未来 dyn* 类型可能原生解决此问题。