Rust: Rust 并发编程:线程、消息传递与共享状态
最后更新:2026-08-26
并发(Concurrency)是程序"同时"做多件事的能力——Rust 通过所有权系统和类型系统在编译时消除数据竞争,让你写出既高效又安全的并发代码。
如果说单线程程序是"一个人在厨房里从头做到尾",那多线程就是"多个厨师同时工作"——有人切菜、有人炒菜、有人装盘。但厨房里人多了就容易出乱子:两个人同时用一把刀、一个人拿了另一个人正在用的配料。Rust 的并发模型就是"带规则的多厨师厨房"——每个厨师有自己的工具,传递食材走专用通道,共享的调味料只能一个人用。
1. 你将学到
std::thread::spawn创建线程与join等待线程结束move闭包将所有权转移到线程中mpsc::channel通道——多生产者单消费者消息传递Mutex<T>互斥锁——保护共享数据的线程安全访问Arc<T>原子引用计数——让Mutex<T>在多线程间共享Send/Synctrait 的概念——Rust 并发安全的基础
2. 多人协作文档编辑系统的故事
(1) 痛苦:没有并发控制的在线文档
Luna 所在的公司正在开发一个在线的协作文档编辑器。早期版本只允许一个人编辑,团队成员抱怨不断:
- Alice 在写方案时,Bob 也在同一段文字上修改——两个人的修改互相覆盖了
- Carol 想看文档的实时更新,但每次都要刷新页面
- Dave 提交了一个修改,但 Luna 完全不知道是谁改的
- 更糟的是,两个编辑者同时修改了同一条内容,文档直接崩了
- 每当 Luna 想加一个"互斥编辑"功能时,代码就变得一团糟——各种全局变量和锁
"如果能给每个编辑者开一个独立线程、修改请求通过通道传递、文档内容用互斥锁保护就好了……"
(2) Rust 并发模型的方案
use std::thread;
use std::sync::mpsc;
use std::sync::{Arc, Mutex};
use std::time::Duration;
fn main() {
// 创建一个共享文档内容(Arc<Mutex<String>>)
let document = Arc::new(Mutex::new(String::from("# 协作文档\n\n")));
// 创建一个通道,用于传输编辑请求
let (tx, rx) = mpsc::channel::<String>();
// 接收线程:持续处理编辑请求
let doc_for_receiver = Arc::clone(&document);
let receiver = thread::spawn(move || {
for edit in rx {
let mut doc = doc_for_receiver.lock().unwrap();
doc.push_str(&edit);
doc.push('\n');
println!("[接收器] 已应用编辑");
}
});
// 模拟编辑者线程
let tx1 = tx.clone();
thread::spawn(move || {
tx1.send("- Alice: 添加了第一章内容".to_string()).unwrap();
});
let tx2 = tx.clone();
thread::spawn(move || {
tx2.send("- Bob: 修改了第二章标题".to_string()).unwrap();
});
// 主线程也发送一条
tx.send("- Luna: 格式调整".to_string()).unwrap();
// 等待所有编辑发送完成
thread::sleep(Duration::from_millis(100));
drop(tx); // 关闭发送端
receiver.join().unwrap();
// 最终文档内容
let final_doc = document.lock().unwrap();
println!("\n=== 最终文档 ===");
println!("{}", *final_doc);
}
Rust 的并发方案清晰分离了职责:每个编辑者是独立线程,编辑请求通过
mpsc::channel传递,文档内容受Mutex保护。Arc让Mutex可以在多个线程间共享。所有权系统确保不会出现数据竞争。
3. 核心概念
(1) Rust 并发编程体系
graph TB
A[Rust 并发编程] --> B[线程管理]
A --> C[消息传递]
A --> D[共享状态]
A --> E[安全保证]
B --> B1["thread::spawn ||"]
B --> B2["join() 等待结束"]
B --> B3["move 闭包转移所有权"]
C --> C1["mpsc::channel"]
C --> C2["Sender / Receiver"]
C --> C3["send / recv"]
D --> D1["Mutex<T> 互斥锁"]
D --> D2["lock() 获取锁"]
D --> D3["Arc<T> 原子引用计数"]
E --> E1["Send: 所有权可跨线程转移"]
E --> E2["Sync: 引用可跨线程共享"]
E --> E3["编译时消除数据竞争"]
(2) 三种并发方式对比
| 特性 | thread::spawn(线程) | mpsc::channel(通道) | Arc<Mutex<T>>(共享状态) |
|---|---|---|---|
| 核心思想 | 独立执行单元 | 消息传递 | 共享内存 + 锁 |
| 数据传递方式 | move 闭包转移所有权 | send/recv 发送消息 | lock 获取内部值 |
| 适用场景 | 独立任务并行执行 | 生产者-消费者模式 | 多个线程访问同一数据 |
| 优点 | 充分利用多核 CPU | 解耦发送者和接收者 | 直接共享任意类型 |
| 缺点 | 线程间通信复杂 | 不适合频繁小数据 | 可能死锁、性能开销 |
| Rust 特色 | 所有权防止悬空指针 | 编译器确保正确使用 | 编译时防止数据竞争 |
(3) Send 与 Sync trait
| Trait | 含义 | 自动实现条件 | 非自动实现类型 |
|---|---|---|---|
| Send | 类型的所有权可以跨线程转移 | 绝大多数类型自动实现 | Rc<T>(非原子引用计数) |
| Sync | 类型的引用可以跨线程共享 | 绝大多数类型自动实现 | RefCell<T>(非原子内部可变性) |
| T: Send + Sync | 类型可以安全地在线程间共享和传递 | Arc<T>、Mutex<T> |
裸指针 *const T、*mut T |
(4) 线程通信方式选型速查
| 通信方式 | 类型 | 数据流向 | 是否需要锁 | 适用场景 |
|---|---|---|---|---|
| 通道 | mpsc::channel |
单向(Sender→Receiver) | 否 | 生产者-消费者模式 |
| 通道(多生产者) | mpsc::channel + tx.clone() |
多→单 | 否 | 多线程汇总结果 |
| 共享状态 | Arc<Mutex<T>> |
双向 | 是(互斥锁) | 多线程读写同一数据 |
| 共享状态(读写) | Arc<RwLock<T>> |
双向 | 是(读写锁) | 读多写少场景 |
| 原子操作 | AtomicUsize 等 |
双向 | 否(硬件级) | 简单计数器/标志位 |
| 屏障 | Barrier |
同步点 | 否 | 多线程汇合点 |
Rust 不靠运行时检查保证线程安全,而是通过
Send和Sync这两个 marker trait 在编译时检查。如果一个类型不是Send,你试图把它传到另一个线程就会编译报错。
4. 并发编程示例
▶ 示例 1:thread::spawn + join——创建和等待线程(难度 ⭐)
// ============================================
// 演示:thread::spawn 创建线程、join 等待结束
// 模拟:厨房里多个厨师同时做不同的菜
// ============================================
use std::thread;
use std::time::Duration;
fn main() {
println!("=== 厨房开工 ===");
// --- 1. 创建三个线程执行不同的任务 ---
let chef1 = thread::spawn(|| {
for i in 1..=3 {
println!("[厨师A] 切菜中... 第{}刀", i);
thread::sleep(Duration::from_millis(50));
}
"A 完成了切菜"
});
let chef2 = thread::spawn(|| {
for i in 1..=3 {
println!("[厨师B] 炒菜中... 第{}步", i);
thread::sleep(Duration::from_millis(40));
}
"B 完成了炒菜"
});
// 主线程也在工作
for i in 1..=3 {
println!("[主厨] 装盘中... 第{}个", i);
thread::sleep(Duration::from_millis(60));
}
// --- 2. join 等待所有线程结束并获取返回值 ---
let result1 = chef1.join().unwrap();
let result2 = chef2.join().unwrap();
println!("\n=== 厨房收工 ===");
println!("厨师A: {}", result1);
println!("厨师B: {}", result2);
println!("所有工作完成!");
}
输出:
=== 厨房开工 ===
[厨师A] 切菜中... 第1刀
[厨师B] 炒菜中... 第1步
[主厨] 装盘中... 第1个
[厨师A] 切菜中... 第2刀
[厨师B] 炒菜中... 第2步
[主厨] 装盘中... 第2个
[厨师A] 切菜中... 第3刀
[厨师B] 炒菜中... 第3步
[主厨] 装盘中... 第3个
=== 厨房收工 ===
厨师A: A 完成了切菜
厨师B: B 完成了炒菜
所有工作完成!
thread::spawn接收一个闭包,在新的操作系统线程中执行。join()阻塞当前线程直到目标线程结束,并返回Result<T>——如果线程 panic,join()会返回Err。线程的执行顺序由操作系统调度决定,每次运行可能不同。
▶ 示例 2:mpsc::channel——消息传递(难度 ⭐⭐)
// ============================================
// 演示:mpsc::channel 多生产者单消费者消息传递
// 模拟:多个编辑者向文档服务器发送修改请求
// ============================================
use std::thread;
use std::sync::mpsc;
use std::time::Duration;
#[derive(Debug)]
enum EditAction {
Insert { user: String, text: String },
Delete { user: String, line: u32 },
Format { user: String, style: String },
}
fn main() {
println!("=== 协作文档编辑器 ===");
// 创建通道:Sender 可以克隆,Receiver 是唯一的
let (tx, rx) = mpsc::channel::<EditAction>();
// --- 1. 启动接收线程(文档服务器)---
let receiver = thread::spawn(move || {
for action in rx {
match &action {
EditAction::Insert { user, text } => {
println!("[服务器] {} 插入了: {}", user, text);
}
EditAction::Delete { user, line } => {
println!("[服务器] {} 删除了第{}行", user, line);
}
EditAction::Format { user, style } => {
println!("[服务器] {} 应用了格式: {}", user, style);
}
}
// 模拟处理时间
thread::sleep(Duration::from_millis(20));
}
println!("[服务器] 通道关闭,停止接收");
});
// --- 2. 创建第一个编辑者线程(Alice)---
let tx1 = tx.clone();
let editor1 = thread::spawn(move || {
tx1.send(EditAction::Insert {
user: "Alice".to_string(),
text: "第一章:Rust 简介".to_string(),
}).unwrap();
thread::sleep(Duration::from_millis(10));
tx1.send(EditAction::Format {
user: "Alice".to_string(),
style: "标题加粗".to_string(),
}).unwrap();
});
// --- 3. 创建第二个编辑者线程(Bob)---
let tx2 = tx.clone();
let editor2 = thread::spawn(move || {
tx2.send(EditAction::Insert {
user: "Bob".to_string(),
text: "Rust 是一门系统编程语言".to_string(),
}).unwrap();
thread::sleep(Duration::from_millis(10));
tx2.send(EditAction::Delete {
user: "Bob".to_string(),
line: 1,
}).unwrap();
});
// --- 4. 主线程也发送一条消息 ---
tx.send(EditAction::Insert {
user: "系统".to_string(),
text: "自动保存文档".to_string(),
}).unwrap();
// 等待编辑者线程结束
editor1.join().unwrap();
editor2.join().unwrap();
// 关闭发送端——所有 Sender 被 drop 后,Receiver 的 for 循环会自动结束
// tx 在这里被 drop(因为 tx 是主线程最后一个 sender)
println!("所有编辑者已完成工作");
receiver.join().unwrap();
println!("=== 编辑会话结束 ===");
}
// 注意:tx 离开作用域时自动 drop,通道关闭
// 如果需要显式关闭,可以用 drop(tx)
输出:
=== 协作文档编辑器 ===
[服务器] 系统 插入了: 自动保存文档
[服务器] Alice 插入了: 第一章:Rust 简介
[服务器] Bob 插入了: Rust 是一门系统编程语言
[服务器] Alice 应用了格式: 标题加粗
[服务器] Bob 删除了第1行
所有编辑者已完成工作
[服务器] 通道关闭,停止接收
=== 编辑会话结束 ===
mpsc代表"多生产者、单消费者"(Multiple Producer, Single Consumer)。Sender可以通过clone复制给多个线程,但Receiver只能有一个。当所有Sender被 drop 后,通道自动关闭,Receiver的迭代器结束。send返回Result——如果接收端已关闭,会返回Err。
▶ 示例 3:Arc<Mutex<T>>——共享状态的线程安全访问(难度 ⭐⭐⭐)
// ============================================
// 演示:Arc<Mutex<T>> 在多个线程间安全共享数据
// 模拟:多个工人同时修改一个共享计数器
// ============================================
use std::thread;
use std::sync::{Arc, Mutex};
use std::time::Duration;
fn main() {
println!("=== 共享计数器演示 ===");
// 用 Arc<Mutex<i32>> 包裹共享数据
let counter = Arc::new(Mutex::new(0i32));
let mut handles = vec![];
// --- 启动 5 个线程,每个线程对计数器加 10 ---
for id in 0..5 {
let counter_clone = Arc::clone(&counter);
let handle = thread::spawn(move || {
for _ in 0..10 {
// lock() 获取互斥锁——如果锁被其他线程持有,当前线程会阻塞等待
let mut num = counter_clone.lock().unwrap();
*num += 1;
println!("[工人{}] 当前计数: {}", id, *num);
// 锁在离开作用域时自动释放
}
println!("[工人{}] 工作完成", id);
});
handles.push(handle);
}
// --- 等待所有线程结束 ---
for handle in handles {
handle.join().unwrap();
}
// --- 读取最终结果 ---
let final_count = counter.lock().unwrap();
println!("\n=== 最终结果 ===");
println!("计数器最终值: {}", *final_count);
println!("预期值: {} (5个线程 × 10次)", 5 * 10);
}
输出:
=== 共享计数器演示 ===
[工人0] 当前计数: 1
[工人0] 当前计数: 2
[工人1] 当前计数: 3
[工人1] 当前计数: 4
[工人0] 当前计数: 5
...(中间输出因线程调度而不同)
[工人4] 当前计数: 50
[工人4] 工作完成
=== 最终结果 ===
计数器最终值: 50
预期值: 50 (5个线程 × 10次)
Mutex<T>提供互斥访问——同一时间只有一个线程能访问内部数据。lock()返回MutexGuard<T>,它实现了Deref和Drop:离开作用域时自动释放锁。Arc<T>(Atomic Reference Counted)是Rc<T>的多线程版本,使用原子操作保证引用计数的线程安全。Arc::clone增加引用计数而不拷贝数据。
▶ 示例 4:线程间协作——生产者-消费者模式(难度 ⭐⭐⭐)
// ============================================
// 生产者-消费者:多生产者 + 单消费者
// ============================================
use std::sync::mpsc;
use std::thread;
use std::time::Duration;
fn main() {
let (tx, rx) = mpsc::channel();
let tx2 = tx.clone();
thread::spawn(move || {
let items = vec!["苹果", "香蕉", "橙子"];
for item in items {
tx.send(format!("生产者1: {}", item)).unwrap();
thread::sleep(Duration::from_millis(100));
}
println!("生产者1 完成");
});
thread::spawn(move || {
let items = vec!["西瓜", "葡萄"];
for item in items {
tx2.send(format!("生产者2: {}", item)).unwrap();
thread::sleep(Duration::from_millis(150));
}
println!("生产者2 完成");
});
println!("=== 消费者接收 ===");
for msg in rx {
println!(" 收到: {}", msg);
}
println!("所有生产者已退出,消费者结束");
}
输出(顺序可能不同):
=== 消费者接收 ===
收到: 生产者1: 苹果
收到: 生产者2: 西瓜
收到: 生产者1: 香蕉
收到: 生产者1: 橙子
收到: 生产者2: 葡萄
生产者1 完成
生产者2 完成
所有生产者已退出,消费者结束
tx.clone()创建多个发送端,实现多生产者模式。当所有发送端(tx和tx2)被 drop 后,rx迭代器自动结束——不需要手动发送"结束"信号。
▶ 示例 5:综合练习——并行数据处理(难度 ⭐⭐⭐)
// ============================================
// 并行计算:多线程分片求和
// ============================================
use std::sync::{Arc, Mutex};
use std::thread;
fn parallel_sum(data: &[i64], num_threads: usize) -> i64 {
let chunk_size = (data.len() + num_threads - 1) / num_threads;
let result = Arc::new(Mutex::new(0i64));
let mut handles = Vec::new();
for i in 0..num_threads {
let chunk = data[i * chunk_size..(i * chunk_size + chunk_size).min(data.len())].to_vec();
let result = Arc::clone(&result);
handles.push(thread::spawn(move || {
let partial: i64 = chunk.iter().sum();
*result.lock().unwrap() += partial;
partial
}));
}
let mut partials = Vec::new();
for handle in handles {
partials.push(handle.join().unwrap());
}
println!("各线程部分和: {:?}", partials);
*result.lock().unwrap()
}
fn main() {
let data: Vec<i64> = (1..=1000).collect();
let sequential_sum: i64 = data.iter().sum();
println!("=== 串行求和 ===");
println!("1 到 1000 的和: {}", sequential_sum);
println!("\n=== 并行求和 (4 线程) ===");
let parallel_result = parallel_sum(&data, 4);
println!("并行求和结果: {}", parallel_result);
assert_eq!(sequential_sum, parallel_result);
let data2: Vec<i64> = (1..=10_000_000).collect();
let start = std::time::Instant::now();
let _ = data2.iter().sum::<i64>();
let seq_time = start.elapsed();
let start = std::time::Instant::now();
let _ = parallel_sum(&data2, 8);
let par_time = start.elapsed();
println!("\n=== 10M 数据性能对比 ===");
println!("串行: {:?}", seq_time);
println!("并行 (8线程): {:?}", par_time);
}
输出:
=== 串行求和 ===
1 到 1000 的和: 500500
=== 并行求和 (4 线程) ===
各线程部分和: [78126, 218874, 109374, 94126]
并行求和结果: 500500
=== 10M 数据性能对比 ===
串行: [时间]
并行 (8线程): [时间]
parallel_sum将数据分片,每个线程计算部分和,通过Arc<Mutex<i64>>累加到总结果。chunk_size向上取整确保所有数据被覆盖。assert_eq!验证并行结果与串行一致。对于大数据集,多线程可显著提升性能。
❓ 常见问题
thread::spawn 创建的线程和主线程是什么关系?mpsc::channel 的 send 和 recv 会阻塞吗?send 通常不阻塞(缓冲区有空位时立即返回),recv 会阻塞直到收到消息。Mutex<T> 和 Arc<Mutex<T>> 必须一起用吗?Mutex 会不会有死锁?Send 和 Sync 有什么区别?Send 是"所有权可以跨线程转移",Sync 是"引用可以跨线程共享"。📖 小节
thread::spawn创建操作系统线程,接收闭包作为执行体;join()等待线程结束并获取结果move闭包 将值的所有权转移到新线程中,避免悬空引用——这是 Rust 所有权系统在并发中的直接应用mpsc::channel提供多生产者单消费者通道,send发送消息,recv接收消息;所有 Sender drop 后通道自动关闭Mutex<T>提供互斥访问,lock()获取锁(阻塞),MutexGuard离开作用域自动释放锁Arc<T>是Rc<T>的线程安全版本,使用原子操作管理引用计数,是Mutex<T>在多线程间共享的桥梁Send/Sync是 Rust 并发安全的 marker trait——Send允许跨线程所有权转移,Sync允许跨线程共享引用,编译器在编译时检查
📝 作业
- 难度 ⭐:写一个程序,创建 3 个线程分别计算 1..=10、11..=20、21..=30 的和,主线程用
join等待所有线程结束后汇总三个部分和并打印最终结果。 - 难度 ⭐⭐:使用
mpsc::channel实现一个"任务分发器"。创建 1 个生产者线程(生成 10 个任务编号 1-10)和 2 个消费者线程(每个消费者从通道接收任务编号并打印"[消费者X] 处理任务 #N")。确保所有任务都被处理完毕。 - 难度 ⭐⭐⭐:实现一个"共享银行账户"系统。使用
Arc<Mutex<f64>>作为共享余额。创建 4 个线程分别模拟存款操作(每个线程随机存入 10-100 元)。主线程等待所有存款线程结束后读取最终余额。额外要求:在每次存款前后打印余额变化,验证没有数据竞争(最终余额 = 初始余额 + 所有存款之和)。