std::thread
std::thread 是标准库中专门负责线程 (thread) 的模块, 直接封装了操作系统提供的线程 API. 在 Rust 里, 每个线程都有自己的栈, 多个线程并发执行, 共享同一份进程内存, 所以需要小心处理数据竞争.
创建线程 (thread::spawn)
thread::spawn 接收一个闭包, 在后台创建一个新线程并立刻执行这个闭包, 返回值是一个 JoinHandle. 新线程的代码和主线程是并发执行的, 两者的输出顺序不确定.
| 方法 | 作用 | 注意点 |
|---|---|---|
| spawn | 创建并启动一个新线程 | 闭包要求 Send + 'static, 返回值是 JoinHandle |
| Builder::new | 创建一个线程构建器 | 可以继续设置线程名、栈大小 |
| Builder::spawn | 按构建器的配置启动线程 | 返回 io::Result<JoinHandle> |
use std::thread;
use std::time::Duration;
fn main() {
// spawn 创建一个新线程, 返回值是 JoinHandle
let handle = thread::spawn(|| {
// current() 获取当前线程, id() 获取线程 id
println!("子线程开始工作, 线程 id = {:?}", thread::current().id());
thread::sleep(Duration::from_millis(200)); // 模拟耗时操作
println!("子线程工作完成");
});
// 主线程继续往下执行, 和子线程是并发的
println!("主线程继续执行, 主线程 id = {:?}", thread::current().id());
// join 会阻塞等待子线程结束 (下面的章节会详细讲)
handle.join().unwrap();
println!("所有线程都已结束");
}
// 主线程和子线程的输出顺序不确定 (并发执行)注意点
thread::spawn 要求传入的闭包满足 'static 生命周期, 也就是说闭包不能借用 main 函数里的局部变量 (后面会讲到用 move 解决)。
等待线程结束 (join)
JoinHandle 是 spawn 返回的"线程句柄", 调用它的 join 方法会阻塞当前线程, 直到对应线程执行完毕, 并且可以取回线程闭包的返回值. 主线程的 main 函数一返回, 整个进程就会直接退出, 不会等待其他线程, 所以想让子线程干完活, 必须手动 join.
| 方法 | 作用 | 注意点 |
|---|---|---|
| JoinHandle::join | 阻塞等待线程结束, 并取回闭包的返回值 | 返回 Result, 线程 panic 时是 Err |
| JoinHandle::is_finished | 判断线程是否已经结束 | 不阻塞, 立即返回 |
| JoinHandle::thread | 拿到该线程的 Thread 句柄 | 可以用来调用 unpark 唤醒 |
use std::thread;
fn main() {
let mut handles = vec![];
// 经典场景: 把一个大任务拆成 4 个小任务, 并行计算
for i in 0..4 {
let handle = thread::spawn(move || {
// 每个线程计算一段等差数列的和
let sum: u32 = (1..=10).map(|x| x + i * 10).sum();
sum // 闭包的最后一行就是线程的返回值
});
handles.push(handle);
}
// join 会阻塞等待线程结束, 并拿到线程的返回值
let mut total = 0;
for handle in handles {
let part = handle.join().unwrap(); // 线程正常结束返回 Ok
println!("某个子任务的结果: {part}"); // 55 / 155 / 255 / 355
total += part;
}
println!("汇总结果: {total}"); // 820 (每次运行一致)
}注意点
子线程在 join 之前就 panic 的话, join 会返回 Err, 直接 unwrap 会让主线程也跟着 panic。业务代码里建议用 match 或 unwrap_or_else 处理。
use std::thread;
use std::time::Duration;
fn main() {
// 反例: 不 join 会怎样?
thread::spawn(|| {
thread::sleep(Duration::from_secs(1));
println!("子线程: 我还没干完活...");
});
// main 函数立刻返回, 整个进程直接退出
// 子线程被进程退出直接终止, 它的 println! 根本没机会执行
println!("main 返回了");
}
// 运行结果只有: main 返回了
// 子线程不会打印任何东西, 因为进程已经退出了注意点
网上常有人说"不 join 也没关系, 进程会等所有线程", 这是错的 (不同语言行为不同)。实测 Rust 的行为是: main 一返回, 进程立即退出, 未结束的线程被直接终止。所以每个 spawn 的句柄都要 join (或改用 thread::scope)。
move 闭包与所有权
thread::spawn 的闭包是 'static 的, 如果闭包要使用外部变量, 必须用 move 关键字把变量的所有权移动进闭包. 不加 move 的话, 闭包会尝试借用局部变量, 编译器直接报错.
| 捕获方式 | 说明 | 注意点 |
|---|---|---|
| FnOnce | 只能调用一次的闭包, 会消耗捕获的变量 | spawn 只要求闭包是 FnOnce |
| FnMut | 可以修改捕获的变量 (可变借用) | 只在单线程内使用 |
| Fn | 只读借用捕获的变量 | 可以被多次调用 |
| move | 强制把捕获的变量移动进闭包内部 | 不加 move 借用局部变量, spawn 会编译报错 |
use std::thread;
fn main() {
// 经典场景: 闭包要使用外部变量, 必须用 move 把所有权移进线程
let data = String::from("需要共享的数据");
let handle = thread::spawn(move || {
println!("子线程拿到了: {data}");
});
// 如果不加 move, 编译器会直接报错:
// error[E0373]: closure may outlive the current function, but it borrows `data`,
// which is owned by the current function
// 因为 spawn 要求闭包是 'static, 不能借用 main 里的局部变量
handle.join().unwrap();
// 另一个经典场景: 把 Vec 里的元素逐个 move 进线程
let names = vec![
String::from("张三"),
String::from("李四"),
String::from("王五"),
];
let mut handles = vec![];
for name in names { // 循环体里逐个 move 出元素
let handle = thread::spawn(move || {
println!("线程处理: {name}");
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
}
// 三个"线程处理"的输出顺序不确定注意点
move 之后, 变量所有权就进了闭包, main 里不能再使用它。如果想"多个线程共享同一个数据", 应该用 Arc (见下文), 而不是把同一个变量 move 进多个线程 (编译也会报错)。
线程配置与信息 (Builder / current / park / sleep)
thread::Builder 可以给线程设置名字和栈大小; thread::current() 可以在任何线程里拿到"我是谁"; park 和 unpark 是一对挂起/唤醒原语.
| 方法 | 作用 | 注意点 |
|---|---|---|
| Builder::name | 设置线程名字 | 参数是字符串 |
| Builder::stack_size | 设置线程栈大小 | 单位是字节 |
| thread::current | 获取当前线程的 Thread 句柄 | 相当于"我是谁" |
| Thread::id | 获取线程 id | 每个线程唯一 |
| Thread::name | 获取线程名字 | 没设置名字时返回 None |
| thread::park | 挂起当前线程 | 需要别的线程调用 unpark 唤醒 |
| Thread::unpark | 唤醒指定线程 | 提前 unpark 也能生效 |
| thread::sleep | 让当前线程睡一段时间 | 参数是 Duration |
| thread::yield_now | 主动让出 CPU | 让其他线程有机会执行 |
use std::thread;
use std::time::Duration;
fn main() {
// 用 Builder 创建带名字、带栈大小的线程
let handle = thread::Builder::new()
.name("worker-1".to_string())
.stack_size(64 * 1024) // 64KB 栈空间
.spawn(|| {
let t = thread::current();
println!("子线程名字: {:?}, id: {:?}", t.name(), t.id());
thread::yield_now(); // 主动让出 CPU
thread::sleep(Duration::from_millis(100));
println!("worker-1 睡醒了");
})
.expect("创建线程失败");
// park: 让主线程挂起, 由另一个线程负责唤醒
let main_thread = thread::current();
let unsparker = thread::spawn(move || {
thread::sleep(Duration::from_millis(100));
println!("准备唤醒主线程");
main_thread.unpark(); // 唤醒主线程
});
println!("主线程进入 park 挂起");
thread::park();
println!("主线程被唤醒了");
handle.join().unwrap();
unsparker.join().unwrap();
}
// 子线程和主线程的输出顺序不确定注意点
park 有个容易踩的坑: 如果 unpark 在 park 之前被调用, park 会立即返回 (相当于提前给了一张"唤醒券")。依赖严格先后顺序的场景别用它。
作用域线程 (thread::scope)
普通 spawn 要求闭包是 'static, 所以只能 move 或者用 Arc. 如果只是想在线程里借用一下局部变量, 可以用 thread::scope 创建作用域线程, 作用域结束时自动等待所有线程完成.
| 方法 | 作用 | 注意点 |
|---|---|---|
| thread::scope | 创建作用域, 内部线程可以借用外部变量 | 返回前自动等待所有线程结束 |
| ScopedJoinHandle::join | 等待作用域内某个线程结束 | 不调用也没关系, scope 会自动等 |
use std::thread;
fn main() {
let data = vec![1, 2, 3, 4, 5];
// 经典场景: 作用域线程直接借用外部变量, 不需要 move 或 Arc
thread::scope(|s| {
for i in 0..3 {
// 重新借用 data: 这样 move 捕获的是引用, 所有权不会转移
let data = &data;
// 循环变量 i 是基本类型, move 会直接复制
s.spawn(move || {
println!("线程 {i} 读到: {data:?}");
});
}
}); // scope 结束前会等待所有子线程完成
// scope 结束后, data 的所有权还在 main 手里, 可以继续用
println!("scope 结束后还能用: {data:?}");
}
// 三个"线程读到"的输出顺序不确定std::sync
std::sync 是标准库的并发原语模块, 包含 Arc (原子引用计数), Mutex / RwLock (锁), Condvar (条件变量), Barrier (屏障), 以及 mpsc (多生产者单消费者通道), 用于安全地跨线程共享数据.
| 类型 | 用途 |
|---|---|
| std::sync::Arc | 跨线程共享同一个数据 (只读) |
| std::sync::Mutex | 互斥锁, 同一时刻只有一个线程能改数据 |
| std::sync::RwLock | 读写锁, 读多写少时比 Mutex 高效 |
| std::sync::Condvar | 条件变量, 让线程等待某个条件成立 |
| std::sync::Barrier | 屏障, 等所有线程到齐再一起放行 |
| std::sync::mpsc | 多生产者单消费者通道, 线程间传消息 |
Arc 原子引用计数
Arc 和 Rc 一样是引用计数指针, 区别在于 Arc 的计数是原子操作, 可以安全地跨线程共享. 多个线程各自持有一份 Arc::clone, 底层数据只有一份, 最后一个引用释放时数据才被回收.
| 方法 | 作用 | 注意点 |
|---|---|---|
| Arc::new | 创建原子引用计数指针 | 数据分配在堆上 |
| Arc::clone | 增加一个引用 | 底层数据只有一份, 只是计数 +1 |
| Arc::strong_count | 查看当前强引用数量 | 常用于调试 |
| Arc::make_mut | 拿到内部数据的可变引用 | 有多个引用时会先复制一份 (写时复制) |
use std::sync::Arc;
use std::thread;
fn main() {
// 经典场景: 多个线程共享同一个不可变数据
let msg = Arc::new(String::from("我是共享数据"));
println!("初始引用计数: {}", Arc::strong_count(&msg)); // 1
let mut handles = vec![];
for i in 0..4 {
// clone 只是增加引用计数, 底层数据只有一份
let msg_clone = Arc::clone(&msg);
let handle = thread::spawn(move || {
println!("线程 {i} 读到: {msg_clone}");
});
handles.push(handle);
}
println!("spawn 后引用计数: {}", Arc::strong_count(&msg)); // 5
for handle in handles {
handle.join().unwrap();
}
println!("线程结束后引用计数: {}", Arc::strong_count(&msg)); // 1
}
// 线程的打印顺序不确定, 但引用计数的输出是确定的注意点
Rc 的计数不是原子操作, 不能跨线程共享。把 Rc move 进 spawn 会直接编译报错: error[E0277]: Rc<String> cannot be sent between threads safely。跨线程共享请用 Arc。
通道 mpsc
mpsc 是 multi-producer, single-consumer (多生产者, 单消费者) 的通道, 是线程间传递消息最常用的方式. 发送端 Sender 可以随意 clone, 接收端 Receiver 只有一个, 数据按发送顺序先进先出.
| 方法 | 作用 | 注意点 |
|---|---|---|
| channel | 创建无界通道 | 返回 (Sender, Receiver) |
| sync_channel | 创建有界通道 | 容量满时 send 会阻塞 |
| Sender::send | 发送一条数据 | 接收端全部关闭时返回 Err |
| Receiver::recv | 阻塞接收一条数据 | 发送端全部关闭时返回 Err |
| Receiver::try_recv | 非阻塞接收 | 没有数据时立刻返回 Err |
| Receiver::iter | 迭代接收 | 发送端全部关闭后迭代结束 |
use std::sync::mpsc;
use std::thread;
fn main() {
// 经典场景: 主线程发任务, 工作线程处理并回传结果
let (task_tx, task_rx) = mpsc::channel::<u32>();
let (result_tx, result_rx) = mpsc::channel::<u32>();
// 工作线程: 接收任务, 算出结果再回传
let worker = thread::spawn(move || {
// task_rx 的 for 循环会阻塞接收, 发送端关闭后循环结束
for task in task_rx {
let result = task * task;
result_tx.send(result).unwrap();
}
println!("工作线程处理完了所有任务");
});
// 主线程发 5 个任务
for i in 1..=5 {
task_tx.send(i).unwrap();
}
// 关闭任务发送端, 工作线程的 for 循环才能结束
drop(task_tx);
// 主线程收集结果
for result in result_rx {
println!("收到结果: {result}");
// 1 / 4 / 9 / 16 / 25
}
worker.join().unwrap();
// try_recv: 非阻塞接收, 没有数据时立刻返回 Err
let (tx, rx) = mpsc::channel::<i32>();
tx.send(42).unwrap();
match rx.try_recv() {
Ok(v) => println!("try_recv 拿到: {v}"), // try_recv 拿到: 42
Err(_) => println!("通道暂时没有数据"),
}
}标准库的通道功能比较基础, 需要"有界 + 复杂调度"的场景可以直接用第三方库 crossbeam-channel:
use std::sync::mpsc;
use std::thread;
fn main() {
// sync_channel: 标准库的有界通道, 容量只有 3
let (tx, rx) = mpsc::sync_channel::<u32>(3);
// 容量满的时候 send 会阻塞, 直到接收方取走数据
let sender = thread::spawn(move || {
for i in 0..5 {
tx.send(i).unwrap();
println!("已发送: {i}");
}
});
// 接收端一边收, 发送端才能继续发
for _ in 0..5 {
let v = rx.recv().unwrap();
println!("收到: {v}");
}
sender.join().unwrap();
}
// "已发送"和"收到"的顺序可能交错, 但一定会全部完成// Cargo.toml 添加依赖:
// [dependencies]
// crossbeam-channel = "0.5"
use crossbeam_channel::bounded;
use std::thread;
fn main() {
// 与 sync_channel 用法几乎一样, 但功能更强 (支持 select!)
let (tx, rx) = bounded::<u32>(3);
let sender = thread::spawn(move || {
for i in 0..5 {
tx.send(i).unwrap();
println!("已发送: {i}");
}
});
for _ in 0..5 {
let v = rx.recv().unwrap();
println!("收到: {v}");
}
sender.join().unwrap();
}RwLock 读写锁
RwLock 允许多个读线程同时读, 但写的时候只有一个线程, 且写的时候不能有读者. 适合"读多写少"的场景, 比如配置信息、缓存数据.
| 方法 | 作用 | 注意点 |
|---|---|---|
| RwLock::new | 创建读写锁 | 内部数据必须有初始值 |
| RwLock::read | 获取读锁 | 多个读锁可以同时存在; 返回 Result |
| RwLock::write | 获取写锁 | 同一时刻只能有一个写锁; 返回 Result |
| RwLock::try_read | 尝试获取读锁 | 拿不到立刻返回 Err, 不阻塞 |
| RwLock::try_write | 尝试获取写锁 | 拿不到立刻返回 Err, 不阻塞 |
use std::sync::{Arc, RwLock};
use std::thread;
fn main() {
// 经典场景: 读多写少的共享数据, 用读写锁
let data = Arc::new(RwLock::new(vec![1, 2, 3]));
let mut handles = vec![];
// 4 个读线程, 可以同时持有读锁
for i in 0..4 {
let data = Arc::clone(&data);
let handle = thread::spawn(move || {
// read() 返回 Result<RwLockReadGuard>
if let Ok(guard) = data.read() {
println!("读线程 {i} 看到: {guard:?}");
}
});
handles.push(handle);
}
// 1 个写线程, 写的时候其他线程都要等待
let data_writer = Arc::clone(&data);
let writer = thread::spawn(move || {
// write() 返回 Result<RwLockWriteGuard>
if let Ok(mut guard) = data_writer.write() {
guard.push(4);
println!("写线程修改了数据");
}
});
handles.push(writer);
for handle in handles {
handle.join().unwrap();
}
// join 保证了写操作一定已经完成
println!("最终数据: {:?}", *data.read().unwrap()); // [1, 2, 3, 4]
}
// 读线程的打印顺序不确定Barrier 与 Condvar
Barrier 是"集合点": 等 N 个线程全部到达后再一起放行, 常用于多线程并行计算的同步阶段. Condvar 是条件变量: 线程可以挂起等待"某个条件成立", 由别的线程 notify 唤醒.
| 方法 | 作用 | 注意点 |
|---|---|---|
| Barrier::new | 创建屏障, n 个线程都到达后才放行 | n 必须和实际线程数一致 |
| Barrier::wait | 等待其他线程到达集合点 | 最后一个到达的线程返回 is_leader() == true |
| Condvar::wait | 释放锁并挂起, 等待通知 | 必须配合 Mutex 使用 |
| Condvar::notify_one | 唤醒一个等待的线程 | 没有等待者时什么都不做 |
| Condvar::notify_all | 唤醒所有等待的线程 | 被唤醒后还要重新检查条件 |
use std::sync::{Arc, Barrier};
use std::thread;
fn main() {
// Barrier: 3 个线程都到达集合点之后, 才一起放行
let barrier = Arc::new(Barrier::new(3));
let mut handles = vec![];
for i in 0..3 {
let barrier = Arc::clone(&barrier);
let handle = thread::spawn(move || {
println!("线程 {i} 到达集合点");
barrier.wait(); // 阻塞等待其他线程
println!("线程 {i} 通过集合点");
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
println!("所有线程都通过了集合点");
}
// "到达集合点" 和 "通过集合点" 各自内部的顺序不确定use std::sync::{Arc, Condvar, Mutex};
use std::thread;
use std::time::Duration;
fn main() {
// Condvar 经典场景: 等待某个条件成立, 而不是傻等
// 注意: Condvar 必须和一个 Mutex 搭配使用
let pair = Arc::new((Mutex::new(0u32), Condvar::new()));
// 等待线程: 等 counter 增加到 5 再继续
let waiter = {
let pair = Arc::clone(&pair);
thread::spawn(move || {
let (lock, cvar) = &*pair;
let mut counter = lock.lock().unwrap();
// wait 会释放锁并挂起, 被 notify 唤醒后再重新拿锁
// 条件不满足就继续等, 所以用 while 而不是 if
while *counter < 5 {
counter = cvar.wait(counter).unwrap();
}
println!("条件满足, counter = {counter}"); // 条件满足, counter = 5
})
};
// 主线程慢慢把 counter 加到 5, 每次加完通知一下
for _ in 0..5 {
let (lock, cvar) = &*pair;
let mut counter = lock.lock().unwrap();
*counter += 1;
cvar.notify_one();
drop(counter); // 尽快释放锁
thread::sleep(Duration::from_millis(50));
}
waiter.join().unwrap();
println!("演示结束");
}注意点
Condvar::wait 被唤醒后并不保证条件已经成立 (可能是"假唤醒"), 所以一定要用 while 循环重新检查条件, 不能只用 if。
std::sync::Mutex
Mutex (互斥锁) 是最常用的锁, 保证同一时刻只有一个线程能访问被保护的数据, 避免数据竞争. 它和 Arc 是绝配: Arc 负责共享, Mutex 负责保护.
基本使用
Mutex::new 创建锁, lock() 返回一个 MutexGuard, 通过守卫才能访问内部数据; 守卫在作用域结束 (或 drop) 时自动解锁, 不需要手动调用 unlock.
| 方法 | 作用 | 注意点 |
|---|---|---|
| Mutex::new | 创建互斥锁 | 内部数据必须提供初始值 |
| Mutex::lock | 获取锁, 返回 MutexGuard | 锁被占用时会阻塞等待; 返回 Result |
| Mutex::try_lock | 尝试获取锁 | 拿不到立刻返回 Err, 不阻塞 |
| Mutex::get_mut | 直接拿到内部数据的可变引用 | 只有独占 Mutex 时才能调用, 不需要上锁 |
| MutexGuard | 锁的"守卫", 通过它访问内部数据 | Drop 时自动解锁 |
use std::sync::{Arc, Mutex};
use std::thread;
fn main() {
// 经典场景: 多个线程往同一个共享列表里追加数据
let list = Arc::new(Mutex::new(Vec::new()));
let mut handles = vec![];
for i in 0..5 {
let list = Arc::clone(&list);
let handle = thread::spawn(move || {
for _ in 0..3 {
// lock() 返回 MutexGuard, 离开作用域时自动解锁
let mut guard = list.lock().unwrap();
guard.push(i);
}
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
let guard = list.lock().unwrap();
println!("列表长度: {}", guard.len()); // 15 (每次运行一致)
println!("内容: {:?}", *guard); // 元素顺序不确定
}锁中毒 (Poisoning)
如果一个线程在持有锁的时候 panic, 锁会进入中毒状态. 之后其他线程再 lock() 会得到 Err(PoisonError), 防止读到"写到一半"的数据. PoisonError 里其实装着被保护的数据, 用 into_inner 可以取出来继续用.
| 方法 | 作用 | 注意点 |
|---|---|---|
| Mutex::lock | 获取锁, 锁中毒时返回 Err(PoisonError) | 直接 unwrap 会 panic |
| PoisonError::into_inner | 取出被锁保护的数据 (守卫) | 相当于"忽略中毒, 继续用" |
| PoisonError::get_ref | 拿到数据的只读引用 | 不转移所有权 |
use std::sync::{Arc, Mutex};
use std::thread;
fn main() {
let lock = Arc::new(Mutex::new(10));
// 这个线程拿到锁之后直接 panic
// 守卫在展开 (unwind) 时没有解锁, 锁进入中毒状态
let bad = {
let lock = Arc::clone(&lock);
thread::spawn(move || {
let _guard = lock.lock().unwrap();
panic!("模拟业务出错"); // 持有锁时 panic
})
};
let _ = bad.join(); // join 返回 Err, 这里直接忽略
// 此时 lock() 返回 Err(PoisonError)
// 如果直接 unwrap 会再次 panic:
// let v = lock.lock().unwrap();
// 正确处理: 从 PoisonError 里把内部的数据取出来
let guard = lock.lock().unwrap_or_else(|e| e.into_inner());
println!("锁中毒了, 但数据还能用: {}", *guard); // 锁中毒了, 但数据还能用: 10
}
// 运行时会先打印子线程的 panic 信息 (输出到 stderr), 主线程正常结束注意点
持有锁的时间越短越好: 不要在上锁的状态里做耗时的 IO, 不然其他线程会一直阻塞等待。守卫在作用域结束时自动解锁, 注意别让守卫活得太久。
经典场景: 多线程计数器
10 个线程, 每个线程对共享计数器加 1000 次, 用 Mutex 保护, 最后结果一定是 10000. 这也是 Arc + Mutex 最经典的组合用法.
use std::sync::{Arc, Mutex};
use std::thread;
fn main() {
// 10 个线程各加 1000 次, 结果必须是 10000
let counter = Arc::new(Mutex::new(0u64));
let mut handles = vec![];
for _ in 0..10 {
let counter = Arc::clone(&counter);
let handle = thread::spawn(move || {
for _ in 0..1000 {
let mut guard = counter.lock().unwrap();
*guard += 1;
}
});
handles.push(handle);
}
// 一定要先 join 所有线程, 再读最终结果
for handle in handles {
handle.join().unwrap();
}
println!("最终计数: {}", *counter.lock().unwrap()); // 10000 (每次运行一致)
}注意点
如果不 join 就直接打印计数, 会读到"加到一半"的值, 因为 main 返回时进程直接退出, 子线程可能还没跑完。
原子类型
原子类型介绍
原子类型 (atomic) 是 CPU 级别的"无锁"原语, 对它的读写是原子的, 不会出现数据竞争, 比 Mutex 更轻量. 标准库提供 AtomicBool / AtomicI32 / AtomicUsize / AtomicU64 等, 分布在 std::sync::atomic 模块.
| 方法 | 作用 | 注意点 |
|---|---|---|
| AtomicUsize::new | 创建原子整数 | 还有 AtomicBool / AtomicI32 / AtomicU64 等 |
| load | 原子地读取当前值 | 需要传内存序 Ordering |
| store | 原子地写入一个新值 | 需要传内存序 Ordering |
| fetch_add | 原子地加一个值, 返回旧值 | 类似 += |
| fetch_sub | 原子地减一个值, 返回旧值 | 类似 -= |
| compare_exchange | CAS 操作: 值等于期望值才写入 | 常用于实现无锁算法 |
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::thread;
fn main() {
// 经典场景: 多线程累加计数, 用原子类型比 Mutex 更轻量
let counter = Arc::new(AtomicUsize::new(0));
let mut handles = vec![];
for _ in 0..10 {
let counter = Arc::clone(&counter);
let handle = thread::spawn(move || {
for _ in 0..1000 {
// fetch_add: 原子地 +1, 并返回旧值
// Relaxed 是最宽松的内存序, 对这个例子足够了
counter.fetch_add(1, Ordering::Relaxed);
}
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
// load: 原子地读取当前值
println!("最终计数: {}", counter.load(Ordering::Relaxed)); // 10000 (每次运行一致)
}注意点
Relaxed 内存序只保证单次操作的原子性, 不保证操作之间的顺序。如果要实现"先写后读"这样的跨线程顺序约定, 需要 Acquire / Release 甚至 SeqCst, 初学者先用 Relaxed 或 SeqCst 即可。
线程安全
Send 与 Sync trait
Send 和 Sync 是 Rust 中两个标记 trait, 编译器用它们保证线程安全: Send 表示所有权可以安全地转移到另一个线程, Sync 表示可以安全地被多个线程同时共享引用.
| trait | 含义 | 注意点 |
|---|---|---|
| Send | 值可以被 move 到另一个线程 | 大多数类型都实现; Rc / 裸指针没有 |
| Sync | 可以被多个线程同时共享引用 | T: Sync 等价于 &T: Send |
| thread::spawn | 要求闭包捕获的内容是 Send + 'static | 不满足就编译报错 |
use std::sync::{Arc, Mutex};
use std::thread;
// 编译期断言: T 必须同时实现 Send 和 Sync
fn assert_send_sync<T: Send + Sync>() {}
fn main() {
// 基本类型、Arc、Mutex 都是 Send + Sync
assert_send_sync::<i32>();
assert_send_sync::<Arc<i32>>();
assert_send_sync::<Mutex<i32>>();
// Rc 既不是 Send 也不是 Sync, 不能跨线程
// assert_send_sync::<std::rc::Rc<i32>>();
// error[E0277]: `Rc<i32>` cannot be sent between threads safely
// spawn 要求闭包捕获的内容满足 Send + 'static
// Arc 满足要求, 可以放心 move 进线程
let shared = Arc::new(42);
let handle = thread::spawn(move || {
println!("子线程读到了: {shared}");
});
handle.join().unwrap();
println!("Send 和 Sync 检查通过");
}相关开源库
标准库的线程和锁已经够用, 但生态里还有更强大、更方便的并发库:
- rayon: 数据并行库, 一行代码把普通迭代器变成多线程并行, 自动管理线程池
- crossbeam: 提供更强大的通道 (select!) 和无锁数据结构
- parking_lot: 比标准库更快的 Mutex / RwLock, 很多知名项目都在用
- tokio: 异步运行时, 用很少的线程处理海量并发连接 (不是传统意义的多线程)
// Cargo.toml 添加依赖:
// [dependencies]
// rayon = "1"
use rayon::prelude::*;
fn main() {
// par_iter: 自动用多线程并行迭代, 用法和 iter 几乎一样
let nums: Vec<u64> = (1..=100).collect();
let sum_of_squares: u64 = nums.par_iter().map(|x| x * x).sum();
println!("1 到 100 的平方和: {sum_of_squares}"); // 338350 (每次运行一致)
// 经典场景: 大数组并行求和
let total: u64 = (1..=10_000_000u64).into_par_iter().sum();
println!("1 到 10000000 的和: {total}"); // 50000005000000 (每次运行一致)
}