课程概览 · 第 16 章

上一篇:智能指针:让资源归属可组合

下一篇:Unsafe 与常用 trait:把不安全性封装在最小边界

线程解决并行执行,异步解决大量等待中的任务。两者都不应从"怎么启动任务"开始,而要先明确:数据归谁、谁能修改、任务何时结束或取消。本章用任务队列业务域把这三条边界逐一落到可运行的代码上——先是标准库的线程与消息传递,再是一个独立的 Tokio 项目演示超时、背压、通道关闭与取消。

线程并行、async await、跨任务所有权、Future 取消与短同步边界的关系
图:线程与 async 的差异在于并行和等待;上线前仍要明确数据归属、取消、关闭与锁作用域。

学习目标与默认选择

学完本章你应当能够:

  • std::thread::scope 写借数据的短命并行计算,并检查 join 结果;
  • std::sync::mpsc 通过所有权转移在 worker 之间传数据;
  • 写一个临界区最小化的 Arc<Mutex<_>> 共享修改示例;
  • Send/Sync 当作编译期属性理解,而不是需要手写的承诺;
  • 在独立的 Tokio 项目里区分四类事件:完成、超时、通道关闭、取消

默认选择:先用单线程把程序写对;需要并行计算时用 std::thread(少量任务)或数据并行库;需要大量 I/O 等待时才引入异步运行时。Arc<Mutex<_>> 是共享可变状态的兜底,不是第一选择——能用通道传递所有权就用通道。

概念讲解:并发与异步解决不同问题

操作系统线程可以在多个 CPU 核上真正并行,但每个线程有独立栈和调度成本。async fn 会被编译为状态机,.await 是状态机可能暂停并把控制权交回运行时的位置;它不会创建线程,也不会让 CPU 密集任务自动并行。运行时只轮询就绪的 Future,因此在 async 任务里执行长时间同步计算会阻塞同一执行器上的其他任务。

一句话分工:CPU 密集 -> 线程;I/O 等待密集 -> async。两者的所有权规则完全相同,这正是本章反复强调的主线。

std::thread::scope:借用数据的并行求和

thread::scope 允许派生线程借用当前函数的数据——作用域结束前所有线程必然 join,借用因此安全:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
use std::thread;

fn sum_parallel(values: &[i64]) -> i64 {
thread::scope(|scope| {
let (left, right) = values.split_at(values.len() / 2);
let left_task = scope.spawn(|| left.iter().sum::<i64>());
let right_sum: i64 = right.iter().sum();
left_task.join().expect("worker panicked") + right_sum
})
}

fn main() {
let values: Vec<i64> = (1..=1000).collect();
assert_eq!(sum_parallel(&values), 500_500);
assert_eq!(sum_parallel(&[]), 0);
assert_eq!(sum_parallel(&[7]), 7);
}

三个要点:spawn 的闭包借用 left,无需 Arcjoin() 的返回值是 Result,panic 会以 Err 形式传回,必须检查;空切片和单元素切片也要被断言覆盖——它们是 split_at 的边界。

mpsc:用所有权转移代替共享

多个 worker 各自生成报告发给汇总方,String 的所有权随 send 转移,任何时刻每个字符串只有一个所有者,不存在共享可变状态:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
use std::sync::mpsc;
use std::thread;

fn main() {
let (tx, rx) = mpsc::channel::<String>();
let mut handles = Vec::new();
for worker_id in 0..3 {
let tx = tx.clone();
handles.push(thread::spawn(move || {
// 所有权随 send 转移到接收方;worker 不再持有该 String
tx.send(format!("worker-{worker_id} 完成")).expect("接收方存活");
}));
}
// 原始 tx 在主线程 drop 后通道才关闭
drop(tx);

let mut reports = Vec::new();
for report in rx {
reports.push(report);
}
for h in handles {
h.join().expect("worker 线程 panic");
}
reports.sort();
assert_eq!(
reports,
vec![
"worker-0 完成".to_string(),
"worker-1 完成".to_string(),
"worker-2 完成".to_string()
]
);
}

注意两个所有权细节:tx.clone() 克隆的是发送端句柄(多个生产者),不是数据;drop(tx) 是必要的一步——只要还有任何发送端存活,for report in rx 就不会结束。

最小合理的 Arc<Mutex<_>>

当多个线程真的要修改同一份统计计数时,才轮到 Arc<Mutex<_>>。示例刻意保持临界区最短(只有内存写),并检查每次 join

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
use std::sync::{Arc, Mutex};
use std::thread;

#[derive(Debug, Default)]
struct Counters {
completed: u64,
failed: u64,
}

fn main() {
let shared = Arc::new(Mutex::new(Counters::default()));
let mut handles = Vec::new();
for _ in 0..8 {
let shared = Arc::clone(&shared);
handles.push(thread::spawn(move || {
// 临界区短:只做内存写,锁在语句末尾立即释放
let mut counters = shared.lock().expect("计数器锁中毒");
counters.completed += 1;
if counters.completed % 4 == 0 {
counters.failed += 1;
}
}));
}
// join 结果必须检查:worker panic 会被 join 返回 Err 暴露出来
for h in handles {
h.join().expect("worker 线程 panic");
}
let counters = shared.lock().expect("计数器锁中毒");
assert_eq!(counters.completed, 8);
assert_eq!(counters.failed, 2);
}

锁的规则:临界区里只放内存操作,不要持锁做 I/O 或调用未知代码;如果某个线程在持锁时 panic,锁会"中毒",后续 lock() 返回 Err——示例用 expect 让这种失败显式爆炸而不是静默错下去。

SendSync:编译期边界,不是需要手写的契约

Send 表示类型的值可以转移到另一个线程;Sync 表示对它的共享引用 &T 可以跨线程共享(等价于 &T: Send)。它们几乎总是由字段自动推导:全 Send 字段的 struct 自动 SendRc<T>!Send(引用计数非原子),把它 move 进线程就会在编译期被拒绝:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
无法编译:Rc 不是 Send,不能跨线程转移(E0277)。

use std::rc::Rc;
use std::thread;

fn main() {
let shared = Rc::new(vec![1, 2, 3]);
thread::spawn(move || {
println!("{:?}", shared); // ERROR: `Rc<Vec<i32>>` cannot be sent
// between threads safely (E0277)
});
}

修正(可运行):单线程共享所有权改用 Arc(原子引用计数,Send + Sync)。

use std::sync::Arc;
use std::thread;

fn main() {
let shared = Arc::new(vec![1, 2, 3]);
let handle = {
let shared = Arc::clone(&shared);
thread::spawn(move || {
assert_eq!(shared.len(), 3);
})
};
handle.join().expect("worker panicked");
assert_eq!(shared.len(), 3);
}

手写 unsafe impl Send/unsafe impl Sync 等于亲自向编译器承诺"没有数据竞争和悬垂访问"——那是不安全代码章节的话题,日常业务代码永远不需要。

概念讲解:async 需要运行时和取消策略

async fn 返回 Future,它不会自动运行,也不会自己创建线程。应用要选择运行时(本章用 Tokio),再用有界通道、超时、取消管理任务的生命周期。这三件事正是 async 最容易出事的地方:

  • 背压:无界队列会把慢消费者变成内存炸弹;有界通道的 send().await 在队列满时让出执行权,压力反向传给生产者。
  • 超时:每个外部等待都应有限时;tokio::time::timeout 把"永远等"变成可观测的错误分支。
  • 取消JoinHandle 被 drop 不会停止任务;真正的取消要么 drop 接收端让对端 send/recv 失败,要么显式 abort()

下面是一个完整独立的 Tokio 项目(课程中唯一的外部依赖示例)。它用一个有界 mpsc 队列喂任务给 awaited worker,用 timeout 切断慢任务,最后分别演示通道关闭、drop 取消与 abort 取消。无网络依赖,四个事件在输出里清晰可辨:

1
2
3
4
5
6
7
8
9
文件:Cargo.toml(完整内容)

[package]
name = "task-workers"
version = "0.1.0"
edition = "2024"

[dependencies]
tokio = { version = "1", features = ["macros", "rt-multi-thread", "sync", "time"] }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
文件:src/main.rs(完整内容)

use std::time::Duration;
use tokio::sync::mpsc;
use tokio::time::timeout;

#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
struct TaskId(u64);

impl TaskId {
fn new(value: u64) -> Result<Self, String> {
if value == 0 {
Err("TaskId 不能为 0".to_string())
} else {
Ok(TaskId(value))
}
}
}

#[derive(Debug, Clone, PartialEq)]
#[allow(dead_code)] // 课程共享词汇:完整枚举在各例中统一出现
enum TaskState {
Todo,
InProgress { started_at: u64 },
Done { finished_at: u64 },
Cancelled { reason: String },
}

#[derive(Debug, Clone, PartialEq)]
struct Task {
id: TaskId,
title: String,
state: TaskState,
labels: Vec<String>,
}

/// 每个 Job 的处理耗时由 worker 决定,用 JobKind 区分演示场景。
#[derive(Debug, Clone, Copy, PartialEq)]
enum JobKind {
/// 立即完成
Fast,
/// 超过超时阈值,用来演示 timeout
Slow,
}

#[derive(Debug, Clone)]
struct Job {
task: Task,
kind: JobKind,
}

fn make_jobs() -> Vec<Job> {
let mut jobs = Vec::new();
for (i, title) in ["索引重建", "缓存预热", "报表导出", "慢备份"].iter().enumerate() {
let id = u64::try_from(i + 1).expect("id fits");
let kind = if *title == "慢备份" { JobKind::Slow } else { JobKind::Fast };
jobs.push(Job {
task: Task {
id: TaskId::new(id).expect("nonzero id"),
title: (*title).to_string(),
state: TaskState::Todo,
labels: Vec::new(),
},
kind,
});
}
jobs
}

#[tokio::main]
async fn main() {
// 有界队列:容量 2。生产者 put 满后会等待,这就是背压。
let (tx, mut rx) = mpsc::channel::<Job>(2);

// 生产者把 4 个 Job 依次入队;队列容量 2 会在第 3 个 send 上等待,
// 直到 worker 取走数据。这正是"有界通道 + await"的背压形态。
let producer = tokio::spawn(async move {
for job in make_jobs() {
// 有界 send:队列满时 .await 让出执行权,而不是无限堆积。
if tx.send(job).await.is_err() {
println!("[producer] 通道已关闭,停止投递");
return;
}
}
// tx 在此被 drop,通道随之关闭
});

let mut completed = 0usize;
let mut timed_out = 0usize;

// worker 循环:recv() 返回 None 表示通道已关闭且排空
while let Some(job) = rx.recv().await {
let title = job.task.title.clone();
let expected = job.kind;
// 单个任务用 timeout 包裹:慢任务不会拖垮整个 worker
let outcome = timeout(Duration::from_millis(50), process(job)).await;
match outcome {
Ok(result) => {
assert_eq!(result, expected);
println!("[worker] 完成: {title}");
completed += 1;
}
Err(_) => {
println!("[worker] 超时: {title} (>50ms)");
timed_out += 1;
}
}
}

// 通道关闭(recv 返回 None,循环正常退出)
println!("[worker] 通道关闭,worker 退出");
producer.await.expect("producer task panicked");

assert_eq!(completed, 3, "3 个快任务应当完成");
assert_eq!(timed_out, 1, "1 个慢任务应当超时");

// drop 取消:未消费的接收端被 drop 后,对端的 send 返回错误,
// 任务可以据此退出。
let (cancel_tx, mut cancel_rx) = mpsc::channel::<Job>(1);
let hung = tokio::spawn(async move {
// 这个任务永远等不到数据,模拟"卡住"的工作
while let Some(_job) = cancel_rx.recv().await {
tokio::time::sleep(Duration::from_secs(600)).await;
}
});
drop(cancel_tx);
// abort:立即取消任务,无论它处于哪个 .await 点。
hung.abort();
match hung.await {
Ok(()) => println!("[cancel] 任务在 abort 前自然结束"),
Err(e) if e.is_cancelled() => println!("[cancel] 任务被 abort 取消"),
Err(e) => panic!("意外的 join 错误: {e:?}"),
}

// 空+关闭的通道:发送端已全部 drop,recv 立即返回 None。
let (tx2, rx2) = mpsc::channel::<Job>(1);
drop(tx2);
let watcher = tokio::spawn(async move {
let mut rx2 = rx2;
assert!(rx2.recv().await.is_none());
});
watcher.await.expect("watcher panicked");
println!("[closed] 空+关闭的通道 recv 返回 None");

println!("全部事件演示完毕: 完成={completed} 超时={timed_out}");
}

async fn process(job: Job) -> JobKind {
match job.kind {
JobKind::Fast => {
tokio::time::sleep(Duration::from_millis(1)).await;
JobKind::Fast
}
JobKind::Slow => {
tokio::time::sleep(Duration::from_millis(500)).await;
JobKind::Slow
}
}
}

运行 cargo run 的输出(时间戳无关,行序确定):

1
2
3
4
5
6
7
8
[worker] 完成: 索引重建
[worker] 完成: 缓存预热
[worker] 完成: 报表导出
[worker] 超时: 慢备份 (>50ms)
[worker] 通道关闭,worker 退出
[cancel] 任务被 abort 取消
[closed] 空+关闭的通道 recv 返回 None
全部事件演示完毕: 完成=3 超时=1

四类事件的判别方式:完成timeout(...) 返回 Ok超时是返回 Err(Elapsed);通道关闭recv() 返回 None 循环退出、或 send() 返回 Err取消是 drop 对端使对方 recv/send 失败、或 abort()JoinHandle::await 返回 JoinError::is_cancelled()

不跨 .await 持锁MutexGuard 不是 Send,把它存活到 .await 之后会被编译器拒绝。这是语言层面的强制,不是风格建议——如果逻辑上必须"锁-等待-再操作",应拆成"锁-取数据-放锁-等待-再锁"。

边界与失败场景

  • join 被忽略spawn 返回的 JoinHandle 被 drop 时 panic 信息会丢失。每个 handle 都要 join 并检查 Err
  • 锁中毒:持锁线程 panic 后 lock() 永远返回 Err;用 expect/unwrap_or_else 显式处理,不要让程序带着坏状态继续跑。
  • 无界队列mpsc::unbounded / tokio::sync::mpsc 无界模式在生产者快于消费者时内存持续增长;默认用有界通道。
  • async 里的同步阻塞:在 .await 之间调用 std::thread::sleep 或重 CPU 计算会卡住整个执行器线程;CPU 工作放 spawn_blocking 或独立线程。
  • 以为 drop JoinHandle 会取消任务:Tokio 里任务继续跑;要停止就 abort() 或关闭它依赖的通道。

为什么可行:所有权规则在线程边界照常生效

Send/Sync 检查发生在 thread::spawn / tokio::spawn编译期:闭包及其捕获值的类型必须证明可以安全离开当前线程。通道之所以能代替锁,是因为 send(value) 把所有权移动过线程边界——接收方拿到独占权,发送方从此无法再碰它,共享可变状态从根上不存在。scope 之所以能借用外部数据,是因为作用域语义保证所有派生线程在借用结束前 join。async 的取消之所以"安全",是因为 .await 点是状态机唯一可能被暂停的位置,future 被 drop 时局部状态随之释放——前提是你没有把外部资源(临时文件、锁、连接)的生命周期挂在 future 之外的某个地方。

常见误区

  • 为了"性能"直接上 Arc<Mutex<...>>:先问能否重排所有权(单一拥有者 + 通道),锁是最后手段。
  • RwLock 一定更快:读写锁有自己的开销和写者饥饿问题;先测量再换。
  • 在 async 任务里做 CPU 密集计算#[tokio::main] 默认多线程运行时,但单个任务阻塞会占住一个 worker 线程;CPU 工作用 spawn_blocking
  • 超时之后忘记资源清理timeout 返回 Err 时被中断的 future 已被 drop,它持有的 RAII 资源会释放;但如果资源在 future 之外注册(如全局表),要手动清理。
  • 测试只跑 happy path:并发 bug 在"通道关闭顺序、重复投递、worker 中途 panic"这些路径上;把它们写进测试矩阵(第 19 章)。

自测

  1. std::thread::scope 的派生线程为什么可以借用栈上的数据,而 thread::spawn 不行?
    答的方向:scope 保证作用域结束前所有线程 join,借用期被限制在安全范围内;spawn 的线程可能活过当前栈帧。
  2. mpsc::channelsend 之后,发送方还能访问那个 String 吗?这消除了哪类 bug?
    答的方向:不能,所有权已转移;消除共享可变状态导致的数据竞争。
  3. SendSync 分别约束什么?Rc<i32> 满足哪个?
    答的方向:Send 约束能否移动到别的线程,Sync 约束 &T 能否跨线程共享;Rc 两者都不满足。
  4. Tokio 中 drop 一个 JoinHandle 会取消任务吗?三种真正取消任务的手段是什么?
    答的方向:不会;abort()、关闭任务依赖的通道(drop 发送端或接收端)、给任务实现协同取消(select/timeout)。
  5. 为什么 MutexGuard.await 会编译失败?如果业务需要在等待期间保持独占,应该怎么改?
    答的方向:guard 非 Send,编译器拒绝;拆成两次加锁,或用异步锁(如 tokio::sync::Mutex)并明确它不是默认选择。