Skip to content

Day 24: 并发基础

现代计算机普遍采用多核处理器,要充分利用硬件能力,程序需要具备并发(concurrency)能力。Rust的标准库提供了强大的并发原语,让我们能够安全地编写多线程程序。

本章将学习Rust中最基础的并发工具:thread模块。我们将讨论如何创建线程、管理线程,以及线程间通信的方法。

使用thread::spawn可以创建新线程:

use std::thread;
use std::time::Duration;
fn main() {
// 创建新线程
let handle = thread::spawn(|| {
for i in 1..5 {
println!("子线程: 第 {} 次迭代", i);
thread::sleep(Duration::from_millis(100));
}
});
// 主线程继续执行
for i in 1..5 {
println!("主线程: 第 {} 次迭代", i);
thread::sleep(Duration::from_millis(100));
}
// 等待子线程完成
handle.join().unwrap();
println!("主线程结束");
}

thread::spawn的参数是一个闭包(closure),这个闭包会在新线程中执行。

join:等待指定线程结束,返回Result。

spawn:创建新线程,立即返回。

use std::thread;
use std::time::Duration;
fn main() {
println!("主线程开始");
// spawn创建线程后不等待其完成
let handle = thread::spawn(|| {
println!("新线程开始");
thread::sleep(Duration::from_millis(500));
println!("新线程结束");
});
// 如果注释掉join(),主线程可能会在新线程结束前结束
// 程序仍会等待所有子线程完成,但这是runtime的行为,不应依赖
println!("主线程继续执行");
// 使用join()明确等待线程结束
handle.join().unwrap();
println!("主线程结束");
}

最佳实践是始终调用join()确保线程完成。如果不join,程序退出时子线程可能被强制终止。

move闭包会捕获环境中的变量所有权,这在线程间传递数据时很重要:

use std::thread;
fn main() {
let data = vec![1, 2, 3];
// move关键字将data的所有权转移到闭包中
let handle = thread::spawn(move || {
println!("线程中访问data: {:?}", data);
// data在这里被使用后,线程结束时会被释放
});
// 这里不能再使用data了,因为所有权已经转移
// println!("主线程访问data: {:?}", data); // 编译错误!
handle.join().unwrap();
}

move是必须的,因为子线程的生命周期可能超过主线程,如果不move,主线程可能会在子线程使用data之前就释放它。

Rust使用消息传递(message passing)进行线程间通信。标准库提供了mpsc(multiple producer, single consumer)channel。

use std::thread;
use std::sync::mpsc;
fn main() {
// 创建channel
let (tx, rx) = mpsc::channel();
// 创建新线程发送消息
thread::spawn(move || {
let msg = String::from("来自新线程的消息");
tx.send(msg).unwrap();
// msg已经被发送,不能再使用
});
// 主线程接收消息
let received = rx.recv().unwrap();
println!("收到消息: {}", received);
}

mpsc表示”多个生产者,单个消费者”。这里我们只用到单个生产者。

多个线程可以向同一个channel发送消息:

use std::thread;
use std::sync::mpsc;
use std::time::Duration;
fn main() {
let (tx, rx) = mpsc::channel();
// 创建多个生产者线程
for i in 0..3 {
let tx_clone = mpsc::Sender::clone(&tx);
thread::spawn(move || {
let msg = format!("线程 {} 的消息", i);
tx_clone.send(msg).unwrap();
});
}
// 手动drop主发送者,否则recv()会一直等待
drop(tx);
// 接收所有消息
for received in rx {
println!("收到: {}", received);
}
}

注意:使用Sender::clone()创建多个发送者,但每个channel只有一个接收者。接收者rx在所有发送者被drop后仍然开放,此时recv()会返回Err。

rx实现了Iterator trait,可以遍历所有消息:

use std::thread;
use std::sync::mpsc;
use std::time::Duration;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
for i in 0..5 {
tx.send(i).unwrap();
thread::sleep(Duration::from_millis(100));
}
});
// 使用for循环接收消息
for received in rx {
println!("收到数字: {}", received);
}
println!("所有消息接收完毕");
}

每个线程都是独立的执行单元。当线程panic时,只会终止自身,不会影响其他线程。

use std::thread;
fn main() {
let handle = thread::spawn(|| {
println!("这个线程会panic");
panic!("故意panic!");
});
// join会返回Result,如果线程panic了会得到Err
let result = handle.join();
match result {
Ok(_) => println!("线程正常结束"),
Err(e) => println!("线程panic了: {:?}", e),
}
// 主线程继续执行,不受子线程panic影响
println!("主线程还在运行");
}

线程panic的设计符合Rust”让错误局部化”的原则:一个线程的问题不会导致整个程序崩溃。

9. thread::scope:更安全的线程管理

Section titled “9. thread::scope:更安全的线程管理”

标准库还提供了thread::scope,它确保所有子线程在scope函数返回前完成:

use std::thread;
fn main() {
// 所有在此scope中spawn的线程都会被join
thread::scope(|scope| {
for i in 0..3 {
scope.spawn(|| {
println!("线程 {}", i);
});
}
// 所有线程保证在这里完成
});
println!("scope已结束,所有线程都已join");
}

scope的好处是自动管理线程生命周期,不需要手动join,也不容易忘记等待线程结束。

今天我们学习了Rust多线程编程的基础:

  1. thread::spawn 创建新线程,接受一个闭包作为线程函数
  2. join 等待线程完成,务必调用以确保线程正确结束
  3. move 闭包将环境变量所有权转移到新线程
  4. mpsc channel 线程间消息传递,发送端可以克隆多个
  5. thread::scope 提供更安全的线程管理,自动join所有子线程

Rust的并发编程强调安全性和明确性。通过所有权和类型系统,Rust能在编译期就消除大部分并发bug,如数据竞争和死锁。