第九章:并发与异步
Rust 的并发编程建立在所有权和类型系统之上,许多并发错误会在编译期被拦截。本章从传统的线程、通道和共享状态开始,逐步深入到 Send/Sync 安全抽象,最后介绍现代异步编程模型 async/await,理解 Rust 如何实现无畏并发。
9.1 线程基础
Rust 标准库提供了 1:1 原生线程支持,每个 Rust 线程对应一个操作系统线程。
创建线程
std::thread::spawn 接收一个闭包,在新线程中执行并返回 JoinHandle:
use std::thread;
use std::time::Duration;
let handle = thread::spawn(|| {
for i in 1..10 {
println!("子线程: {i}");
thread::sleep(Duration::from_millis(1));
}
});
for i in 1..5 {
println!("主线程: {i}");
thread::sleep(Duration::from_millis(1));
}
当主线程结束时,整个程序退出,无论子线程是否完成。因此通常需要用 handle.join() 等待子线程结束。
join 等待线程
join() 会阻塞当前线程,直到被等待的线程完成,返回该线程的 Result(用于处理 panic):
handle.join().unwrap();
move 闭包
闭包默认会从环境中借用值,但在跨线程场景下,必须强制闭包获得所有权,以避免悬垂引用。move 关键字会强制闭包将所用到的环境变量移入其中:
let v = vec![1, 2, 3];
let handle = thread::spawn(move || {
println!("{:?}", v); // v 所有权转移到闭包
});
线程局部存储
thread_local! 宏可以定义线程局部的静态变量,每个线程拥有独立的副本:
use std::cell::RefCell;
thread_local! {
static FOO: RefCell<u32> = RefCell::new(0);
}
FOO.with(|f| {
*f.borrow_mut() = 42;
});
这类变量特别适合存储线程专属的上下文或缓存。
9.2 消息传递(Channel)
Rust 标准库提供了多生产者单消费者(mpsc)通道,通过消息传递在线程间通信,符合“不要通过共享内存来通信,而要通过通信来共享内存”的理念。
创建通道
use std::sync::mpsc;
let (tx, rx) = mpsc::channel(); // tx: 发送端, rx: 接收端
tx.send(value):发送值,若接收端已抛弃则返回Err。rx.recv():阻塞直到收到值,通道关闭时返回Err。rx.try_recv():非阻塞,有值时返回Ok,否则返回Err。
thread::spawn(move || {
tx.send("hello from thread").unwrap();
});
let msg = rx.recv().unwrap();
println!("收到: {msg}");
多生产者
mpsc 的 tx 可以克隆,产生多个发送端共享同一接收端:
let (tx, rx) = mpsc::channel();
let tx1 = tx.clone();
thread::spawn(move || { tx.send("线程1").unwrap(); });
thread::spawn(move || { tx1.send("线程2").unwrap(); });
for received in rx {
println!("{received}");
}
接收端 rx 可作为迭代器使用,阻塞地获取每条消息,直到所有发送端都被丢弃且通道为空时迭代结束。
9.3 共享状态并发
消息传递并非万能,有时需要多个线程同时读写同一数据。Rust 的类型系统让共享内存也安全。
Mutex<T>:互斥锁
Mutex<T> 提供内部数据的互斥访问。lock() 方法返回 LockResult<MutexGuard<T>>,MutexGuard 实现 Deref 和 DerefMut,可以透明地访问内部值,离开作用域时自动释放锁:
use std::sync::Mutex;
let m = Mutex::new(5);
{
let mut num = m.lock().unwrap();
*num = 6;
} // 锁在此处自动释放
println!("{:?}", m);
lock() 在获得锁的线程 panic 时会返回 Err,用 unwrap() 处理是常见做法。
Arc<T> 多线程所有权共享
由于线程之间可能没有明确的所有权层次,需要引用计数来共享所有权。Arc<T>(原子引用计数)功能与 Rc<T> 相同,但通过原子操作保证线程安全:
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!("结果: {}", *counter.lock().unwrap()); // 10
Arc<Mutex<T>> 是多线程共享可变状态最经典的模式,但要注意死锁风险和锁的开销。
RwLock<T>:读写锁
RwLock<T> 允许多个读或一个写,适合读多写少场景:
use std::sync::RwLock;
let lock = RwLock::new(5);
{
let r1 = lock.read().unwrap();
let r2 = lock.read().unwrap(); // 多个读可以共存
}
{
let mut w = lock.write().unwrap();
*w += 1;
}
内部原理上,RwLock 在读多时优于 Mutex,但可能增加锁竞争开销。
原子类型
std::sync::atomic 模块提供无需锁的原子操作,适用于简单的共享标志或计数器:
use std::sync::atomic::{AtomicBool, AtomicI32, Ordering};
use std::sync::Arc;
use std::thread;
let running = Arc::new(AtomicBool::new(true));
let r = running.clone();
thread::spawn(move || {
while r.load(Ordering::SeqCst) {
// 执行工作
}
});
running.store(false, Ordering::SeqCst);
Ordering(Relaxed、Acquire、Release、SeqCst)控制内存顺序保证,SeqCst 是最严格、最容易理解的,但在性能敏感场景下可选用较宽松的顺序。
其他同步原语
Barrier:让所有线程在某点同步,直到达到指定数量后一起继续。Condvar:条件变量,配合Mutex实现线程通知等待。
这些原语为更复杂的同步模式提供了底层支持。
9.4 Send 与 Sync trait
这两个标记 trait 是 Rust 并发安全的基石,编译器会自动为符合规则的类型实现它们。
Send
Send 标记表示该类型的所有权可以安全地在线程间转移。几乎所有 Rust 类型都是 Send,但 Rc<T> 例外,因为它的引用计数不是原子操作。
fn is_send<T: Send>() {}
is_send::<i32>(); // 通过
// is_send::<Rc<i32>>(); // 编译错误
Sync
Sync 标记表示该类型的共享引用 &T 可以安全地在多个线程间访问。即如果 &T 是 Send,则 T 是 Sync。
fn is_sync<T: Sync>() {}
is_sync::<i32>(); // 通过
// is_sync::<RefCell<i32>>(); // 编译错误
Mutex<T> 是 Sync,而 RefCell<T> 不是,因为 RefCell 的借用检查是运行时的且非线程安全。
自动实现与手动 unsafe
这两个 trait 是自动推导的:如果一个类型的所有成员都是 Send / Sync,那该类型也自动是。如果某个类型内部使用了原始指针等需要手动保证线程安全的场景,就需要 unsafe impl Send 等方式来自行承诺,这必须极其谨慎。
9.5 异步基础(async/await)
Rust 的异步编程通过 async/await 语法和 Future trait 实现,没有内置运行时——你需要选择外部执行器(如 tokio、async-std)来驱动任务。
Future trait
异步操作的核心抽象是 Future,它代表一个尚未完成的计算。简化定义如下:
pub trait Future {
type Output;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}
Poll::Pending 表示未就绪,执行器会在就绪时再次调用 poll;Poll::Ready(value) 则返回最终值。
async 块与函数
用 async 关键字标记函数或块,编译器会将其转换为实现了 Future 的状态机:
async fn fetch_data() -> String {
// 异步操作
"data".to_string()
}
let future = fetch_data(); // 返回 Future,尚未执行
执行器与 .await
执行器(如 tokio)负责调度和推进 Future。在 async 函数中使用 .await 等待某个 Future 完成并获取结果:
use tokio;
#[tokio::main] // 使用 tokio 执行器的宏,隐式启动运行时
async fn main() {
let data = fetch_data().await;
println!("{data}");
}
.await 会在该点让出执行权,直到 Future 就绪。
async move 块
与普通闭包一样,async 块默认会从环境中借用;使用 async move 让块获得所用变量的所有权,常用于生成不同生命周期的任务:
let s = String::from("hello");
let task = async move {
println!("{s}");
};
异步任务:tokio::spawn
tokio::spawn 将一个 Future 提交为独立任务并立即在后台运行(需满足 'static 生命周期):
let handle = tokio::spawn(async {
"return value"
});
let result = handle.await.unwrap();
任务失败或 panic 时可以通过 JoinHandle 捕获。
异步 I/O
标准库的同步 I/O 会阻塞线程,异步版本(如 tokio::net、tokio::fs)则配合 .await 实现非阻塞操作:
use tokio::net::TcpListener;
let listener = TcpListener::bind("127.0.0.1:8080").await.unwrap();
loop {
let (socket, _) = listener.accept().await.unwrap();
tokio::spawn(async move {
// 处理 socket
});
}
异步 trait 与 async_trait
自 Rust 1.75(2023 年底)起,稳定版已支持在 trait 中直接定义 async fn(静态分发):
trait Fetcher {
async fn fetch(&self) -> String;
}
但该特性目前仍有限制:这类 trait 不是 dyn 安全的(无法用于 Box<dyn Fetcher> 等动态分发场景),且公共 trait 中的 async fn 默认无法显式声明返回的 Future 是否 Send(可配合 trait-variant crate 的 #[trait_variant::make] 生成 Send 版本)。
若需要动态分发(dyn Trait),仍需使用 async_trait 宏,它将异步方法转换为返回 Pin<Box<dyn Future<Output = ...> + Send>> 的形式:
use async_trait::async_trait;
#[async_trait]
trait Fetcher {
async fn fetch(&self) -> String;
}
Stream:异步迭代器
Stream 类似于 Iterator,但 next 方法返回 Poll<Option<Item>>。futures crate 提供了 StreamExt 扩展,可在流上使用 map、filter、for_each 等组合子:
use futures::stream::{self, StreamExt};
let stream = stream::iter(vec![1, 2, 3]);
stream.for_each(|x| async move {
println!("{x}");
}).await;
常见异步模式
select!(tokio::select!/futures::select!):同时等待多个异步操作,任一完成即返回,可用于超时等。join!(tokio::join!/futures::join!):并行等待多个Future全部完成。- 并发请求:将多个
Future提交到tokio::spawn或收集为Vec后用join_all并发执行。
这些模式使异步代码能以接近同步代码的简洁性实现高并发。
传统并发模型通过线程、通道、互斥锁和原子类型提供了可靠的基础,而 async/await 则让高并发 I/O 密集任务变得轻盈直观。Rust 将安全性从单线程延续到并发世界,Send/Sync 是这一承诺的守护者。掌握这些工具后,你就可以构建既安全又高性能的并发程序了。
参考文献
| 资料 | 说明 |
|---|---|
| Concurrency | 官方书第 16 章 |
| Async | async/await |
| Tokio | 常用异步运行时(第三方 crate) |
相关文章
第八章:测试、文档与质量
Rust 将测试和文档视为语言的一等公民。通过内置的测试框架,你可以在项目里编写单元测试、集成测试以及随文档一起运行的示例测试,三者共用 cargo test 命令。配合 cargo doc 生成文档,形成了一套确保代码质量与可维护性…
泛型与 Trait
泛型和 trait 是 Rust 实现代码复用与多态的两大支柱。前置:Rust 基础。泛型让代码可以工作在多种类型上而不牺牲性能,trait 定义了类型间的共享行为,二者结合形成了零成本抽象的强大表达能力。本章将深入泛型定义、trai…
第二章:基础语法与类型
Rust 的类型系统和语法设计处处体现着“安全”与“显式”的理念。这一章你将掌握变量绑定、基本类型、复合类型、函数定义以及所有基础控制流结构,它们是你写出任何 Rust 程序的基石。
第七章:模块系统与包管理
Rust 的模块系统为代码组织、封装和复用提供了一套严谨但灵活的机制。包、crate、模块以及 use 路径相互配合,让你能够把项目拆解成清晰的功能单元,同时精确控制哪些对外可见。本章将带你系统掌握这些构建大型 Rust 项目所必需的…
第十章:进阶特性与模式
Rust 的核心安全保证覆盖了绝大多数日常编程场景。但当你需要打破常规——无论是编写极致通用的抽象、与 C 库交互,还是内联优化——本章将带你进入 Rust 的深层能力:声明宏与过程宏、unsafe 的超能力与封装、高级类型系统技巧…
第五章:错误处理
Rust 将错误明确分为两类:不可恢复的错误与可恢复的错误。通过 panic! 处理前一种,Result 处理后一种,这让程序的错误路径不再是隐式的控制流,而是强类型、必须处理的代码分支。配合 Option 对缺失值的处理以及丰富的组…
Series
rust
1 / 10