FORMA

第九章:并发与异步

Rust 的并发编程建立在所有权和类型系统之上,许多并发错误会在编译期被拦截。本章从传统的线程、通道和共享状态开始,逐步深入到 Send/Sync 安全抽象,最后介绍现代异步编程模型 async/await,理解 Rust 如何实现无畏并发。

9.1 线程基础

Rust 标准库提供了 1:1 原生线程支持,每个 Rust 线程对应一个操作系统线程。

创建线程

std::thread::spawn 接收一个闭包,在新线程中执行并返回 JoinHandle

rust
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):

rust
handle.join().unwrap();

move 闭包

闭包默认会从环境中借用值,但在跨线程场景下,必须强制闭包获得所有权,以避免悬垂引用。move 关键字会强制闭包将所用到的环境变量移入其中:

rust
let v = vec![1, 2, 3];
let handle = thread::spawn(move || {
    println!("{:?}", v); // v 所有权转移到闭包
});

线程局部存储

thread_local! 宏可以定义线程局部的静态变量,每个线程拥有独立的副本:

rust
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)通道,通过消息传递在线程间通信,符合“不要通过共享内存来通信,而要通过通信来共享内存”的理念。

创建通道

rust
use std::sync::mpsc;

let (tx, rx) = mpsc::channel(); // tx: 发送端, rx: 接收端
  • tx.send(value):发送值,若接收端已抛弃则返回 Err
  • rx.recv():阻塞直到收到值,通道关闭时返回 Err
  • rx.try_recv():非阻塞,有值时返回 Ok,否则返回 Err
rust
thread::spawn(move || {
    tx.send("hello from thread").unwrap();
});

let msg = rx.recv().unwrap();
println!("收到: {msg}");

多生产者

mpsctx 可以克隆,产生多个发送端共享同一接收端:

rust
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 实现 DerefDerefMut,可以透明地访问内部值,离开作用域时自动释放锁:

rust
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> 相同,但通过原子操作保证线程安全:

rust
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> 允许多个读或一个写,适合读多写少场景:

rust
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 模块提供无需锁的原子操作,适用于简单的共享标志或计数器:

rust
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);

OrderingRelaxedAcquireReleaseSeqCst)控制内存顺序保证,SeqCst 是最严格、最容易理解的,但在性能敏感场景下可选用较宽松的顺序。

其他同步原语

  • Barrier:让所有线程在某点同步,直到达到指定数量后一起继续。
  • Condvar:条件变量,配合 Mutex 实现线程通知等待。

这些原语为更复杂的同步模式提供了底层支持。

9.4 SendSync trait

这两个标记 trait 是 Rust 并发安全的基石,编译器会自动为符合规则的类型实现它们。

Send

Send 标记表示该类型的所有权可以安全地在线程间转移。几乎所有 Rust 类型都是 Send,但 Rc<T> 例外,因为它的引用计数不是原子操作。

rust
fn is_send<T: Send>() {}
is_send::<i32>();      // 通过
// is_send::<Rc<i32>>(); // 编译错误

Sync

Sync 标记表示该类型的共享引用 &T 可以安全地在多个线程间访问。即如果 &TSend,则 TSync

rust
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 实现,没有内置运行时——你需要选择外部执行器(如 tokioasync-std)来驱动任务。

Future trait

异步操作的核心抽象是 Future,它代表一个尚未完成的计算。简化定义如下:

rust
pub trait Future {
    type Output;
    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}

Poll::Pending 表示未就绪,执行器会在就绪时再次调用 pollPoll::Ready(value) 则返回最终值。

async 块与函数

async 关键字标记函数或块,编译器会将其转换为实现了 Future 的状态机:

rust
async fn fetch_data() -> String {
    // 异步操作
    "data".to_string()
}

let future = fetch_data(); // 返回 Future,尚未执行

执行器与 .await

执行器(如 tokio)负责调度和推进 Future。在 async 函数中使用 .await 等待某个 Future 完成并获取结果:

rust
use tokio;

#[tokio::main] // 使用 tokio 执行器的宏,隐式启动运行时
async fn main() {
    let data = fetch_data().await;
    println!("{data}");
}

.await 会在该点让出执行权,直到 Future 就绪。

async move

与普通闭包一样,async 块默认会从环境中借用;使用 async move 让块获得所用变量的所有权,常用于生成不同生命周期的任务:

rust
let s = String::from("hello");
let task = async move {
    println!("{s}");
};

异步任务:tokio::spawn

tokio::spawn 将一个 Future 提交为独立任务并立即在后台运行(需满足 'static 生命周期):

rust
let handle = tokio::spawn(async {
    "return value"
});
let result = handle.await.unwrap();

任务失败或 panic 时可以通过 JoinHandle 捕获。

异步 I/O

标准库的同步 I/O 会阻塞线程,异步版本(如 tokio::nettokio::fs)则配合 .await 实现非阻塞操作:

rust
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(静态分发):

rust
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>> 的形式:

rust
use async_trait::async_trait;

#[async_trait]
trait Fetcher {
    async fn fetch(&self) -> String;
}

Stream:异步迭代器

Stream 类似于 Iterator,但 next 方法返回 Poll<Option<Item>>futures crate 提供了 StreamExt 扩展,可在流上使用 mapfilterfor_each 等组合子:

rust
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 章
Asyncasync/await
Tokio常用异步运行时(第三方 crate)

相关文章

Series

rust

1 / 10

Rust 基础