Agent X-Ray
RuntimeNotesAbout
Notes/代码工程/Rust 深度教材/第9章

第9章:并发编程 —— 线程、消息、锁、原子与线程安全

12 分钟 · 更新于 2026-09-01

第9章:并发编程 —— 线程、消息、锁、原子与线程安全

并发程序的难点很少是“如何创建线程”,更多是如何界定任务寿命、数据归属、资源上限和退出协议。本章围绕一个批量缩略图调度器展开:作业通过有界通道进入工作线程,结果通过另一条通道返回,汇总数据由锁保护,轻量指标使用原子变量,最后由 SendSync 解释类型为何能够跨线程流动。

本章目标

  • 为每个线程明确启动、停止、join 和错误传播责任。
  • 使用消息传递表达所有权转移,而不是默认共享全部状态。
  • 理解 ArcMutex 解决的是两类不同问题。
  • 根据状态是否存在复合不变量选择 Atomic 或锁。
  • 能识别无界队列、长时间持锁、死锁和失联线程等活性风险。
  • Send / Sync 约束反推类型设计,而不是绕过编译器。

一、先确定数据怎样流动

并发设计可以从三种数据关系入手:

关系推荐表达例子
一方完成后交给另一方channel作业、事件、处理结果
多线程需要查看或修改同一份状态Arc<Mutex<T>> / Arc<RwLock<T>>缓存、配置、汇总状态
一个独立机器字状态需要低成本更新Atomic计数器、序号、停止标志

缩略图调度器的数据路径如下:

text
提交者 --Job--> 固定 worker 集合 --JobResult--> 汇总者
                     |
                     +-- Arc<Mutex<Summary>>
                     +-- Arc<AtomicUsize>

这里没有哪种原语“更高级”。选择标准是语义:

  • 作业只能由一个 worker 处理,因此移动所有权。
  • 多项统计需要一起保持一致,因此放在同一把锁后。
  • 当前正在执行的数量只是独立指标,可以用 Atomic。

二、线程的生命周期:spawn、move、join

2.1 调度顺序不是程序契约

rust
use std::thread;

let handle = thread::spawn(|| {
    println!("worker started");
});

println!("main continues");
handle.join().unwrap();

两条输出的先后顺序不确定。操作系统可以在任何适当时刻调度线程,因此正确性不能依赖观察到的偶然顺序。

需要先后关系时,应使用 join、channel、锁、barrier 等同步手段明确表达。

2.2 move 解决闭包捕获的寿命问题

rust
use std::thread;

let file_name = String::from("cover.png");
let handle = thread::spawn(move || {
    println!("processing {file_name}");
});

handle.join().expect("worker panicked");

新线程可能比创建它的栈帧活得更久。move 让闭包取得 file_name 的所有权,使线程不再依赖外部局部变量。

move 不承诺深拷贝。对 String 是移动,对 u64Copy 类型则是复制值。

2.3 JoinHandle 是结果和故障边界

rust
use std::thread;

let handle = thread::spawn(|| -> usize {
    6 * 7
});

let answer = handle.join().expect("calculation thread panicked");
assert_eq!(answer, 42);

join() 同时承担:

  • 等待线程完成;
  • 获取线程返回值;
  • 观察线程是否 panic

如果把 JoinHandle 丢弃,线程会继续运行,但调用者失去等待和观察它的直接手段。

启动线程时同步设计退出 每次 spawn 都应该回答:谁发出停止信号、线程怎样离开循环、谁保存句柄、谁执行 join、线程失败后系统如何处理。

2.4 作用域线程适合短期借用

标准库的 scoped thread 允许线程借用当前作用域中的数据,同时保证离开作用域前所有线程都已结束:

rust
use std::thread;

let values = vec![2, 4, 6, 8];

thread::scope(|scope| {
    scope.spawn(|| {
        println!("sum = {}", values.iter().sum::<i32>());
    });
});

这比为了满足 'static 而无条件克隆数据更准确。长期后台线程仍应拥有自己的状态和退出协议。


三、消息传递:让所有权沿协议移动

3.1 最小作业队列

rust
use std::sync::mpsc;
use std::thread;

let (tx, rx) = mpsc::channel::<String>();

let worker = thread::spawn(move || {
    while let Ok(path) = rx.recv() {
        println!("resize {path}");
    }
});

tx.send(String::from("photo-a.jpg")).unwrap();
tx.send(String::from("photo-b.jpg")).unwrap();
drop(tx);

worker.join().unwrap();

send(path) 后,String 的所有权进入通道;接收端拿到消息后成为新所有者。所有发送端销毁后,recv 返回断开错误,循环自然结束。

3.2 用枚举描述协议,而不是魔法值

rust
#[derive(Debug)]
struct ResizeJob {
    id: u64,
    source: String,
    width: u32,
}

#[derive(Debug)]
enum WorkerCommand {
    Process(ResizeJob),
    Stop,
}

枚举可以让协议具有明确分支,并由编译器检查是否处理完整。以后增加 ReloadConfigFlushHealthCheck 时,也无需挤占业务字段。

3.3 多发送者与“永不结束”的接收循环

rust
use std::sync::mpsc;

let (tx, rx) = mpsc::channel();
let another_tx = tx.clone();

std::thread::spawn(move || {
    another_tx.send(10).unwrap();
});

drop(tx);
for value in rx {
    println!("{value}");
}

只要还有任何 Sender 存活,接收者就必须认为未来仍可能出现消息。接收循环迟迟不结束时,第一项排查应是寻找遗留的发送端 clone。

3.4 背压是容量治理

mpsc::channel() 是无界通道。如果生产速度长期高于消费速度,消息会持续占用内存。

rust
use std::sync::mpsc;

let (tx, rx) = mpsc::sync_channel::<ResizeJob>(32);

sync_channel 的缓冲区满后,send 会等待:

  • 容量 0:发送与接收直接会合;
  • 容量 N:最多积压 N 条;
  • 无界:吞吐峰值更平滑,但没有内存上限。

容量不是随意的魔法数字。应结合单条消息大小、可接受排队时间、消费速率和过载策略确定。


四、共享状态:Arc 负责寿命,Mutex 负责访问

4.1 为什么两个类型缺一不可

rust
use std::sync::{Arc, Mutex};

let summary = Arc::new(Mutex::new(Vec::<String>::new()));
  • Arc 允许多个线程共同拥有 Mutex
  • Mutex 允许同一时刻只有一个线程修改内部 Vec

只使用 Arc<Vec<_>> 没有可变访问机制;只使用 Mutex<Vec<_>> 又无法把所有权安全地复制给多个长期线程。

4.2 短临界区示例

rust
use std::sync::{Arc, Mutex};
use std::thread;

let completed = Arc::new(Mutex::new(Vec::<u64>::new()));
let mut handles = Vec::new();

for id in 0..4 {
    let completed = Arc::clone(&completed);
    handles.push(thread::spawn(move || {
        let output = id * id; // 锁外计算

        {
            let mut guard = completed.lock().unwrap();
            guard.push(output);
        } // guard 在此释放
    }));
}

for handle in handles {
    handle.join().unwrap();
}

临界区中只保留真正需要互斥的更新。文件读写、网络访问、休眠、压缩和未知回调都不应在持锁期间执行。

4.3 Poisoning 是风险提示

线程持有 MutexGuard 时发生 panic,标准库会把锁标记为 poisoned。后续 lock() 返回 PoisonError,提醒共享状态可能处于只更新了一半的状态。

应用需要明确恢复策略:

  • 状态不可恢复:记录并终止该子系统;
  • 状态可校验和重建:取回 guard 后修复;
  • 独立任务失败不影响整体:把可变状态更新设计成事务式操作。

不应把所有 lock().unwrap() 都当成永远安全的模板。教学示例可以简化,生产库应向上返回有意义的错误。

4.4 RwLock 不是自动优化

RwLock<T> 允许多个读者或一个写者。它适合读远多于写、读临界区有一定长度的情况。但实现成本、写者饥饿策略和平台差异都会影响效果,是否优于 Mutex 应通过负载测试判断。


五、Atomic:只在状态足够窄时使用

5.1 独立计数

rust
use std::sync::atomic::{AtomicUsize, Ordering};

static FINISHED: AtomicUsize = AtomicUsize::new(0);

fn record_finished() {
    FINISHED.fetch_add(1, Ordering::Relaxed);
}

当计数仅用于指标展示,并不用于证明其他内存已经更新完成时,Relaxed 可以满足“操作本身原子化”的需求。

5.2 内存顺序表达可见性关系

Ordering关注点典型用途
Relaxed仅保证本次原子操作不可分割独立计数、采样指标
Acquire读取后能看到配对发布之前的写入消费就绪状态
Release写入前的操作先于状态发布发布数据已经准备好
AcqRel原子读改写同时承担获取与发布状态机转换
SeqCst提供最强的全局顺序直觉保守实现、验证初期

内存顺序不是简单的性能档位。错误使用 Release/Acquire 可能让代码在测试环境长期正常,却在特定架构或优化下暴露问题。业务代码若不需要无锁协议,应优先选择 channel 或锁。

5.3 复合不变量应集中到锁中

假设系统要求:

text
succeeded + failed == accepted

如果三个字段分别使用 Atomic,单独读取时可能来自不同时间点,无法构成一致快照。将它们放在一个结构体中并用一把 Mutex 保护,更能直接表达约束:

rust
#[derive(Debug, Default, Clone)]
struct Summary {
    accepted: usize,
    succeeded: usize,
    failed: usize,
}

选择准则 单个数值、标志或指针状态可以评估 Atomic;多个字段需要共同更新或共同读取时,优先使用锁保护聚合结构。


六、Send 与 Sync:类型层面的跨线程许可

6.1 两个标记特征

  • Send:某个值的所有权可以安全转移到另一线程。
  • Sync:某个类型的共享引用 &T 可以安全地被多个线程使用。

普通结构体会根据字段自动推导这些性质。例如,结构体中含有 Rc<T>,通常就无法发送到另一个线程。

6.2 常见边界

类型线程语义
Rc<T>引用计数非原子,不是跨线程共享工具
RefCell<T>借用状态无同步保护,不适合线程共享
Cell<T>可通过共享引用写入,不是 Sync
Arc<T>引用计数线程安全,但内部 T 仍需满足约束
Mutex<T>通过锁提供线程间的可变访问
裸指针编译器不自动替作者证明线程安全

因此,Arc<T> 不是给任意 T 加上“线程安全认证”。它只解决引用计数自身的数据竞争问题。

6.3 不要用 unsafe impl 消除报错

rust
// unsafe impl Send for LocalCache {}
// unsafe impl Sync for LocalCache {}

手动实现意味着作者承诺:线程迁移、并发共享、析构、内部别名和所有访问路径都满足 Rust 的内存模型。除非正在实现底层同步容器,并且能够给出完整证明,否则应修改字段或线程边界,而不是强行标记。


七、实战:固定 worker 的缩略图调度器

标准库的 mpsc::Receiver 只有一个接收者。为了让多个 worker 共同获取任务,可以把接收端放入 Arc<Mutex<_>>。这会让“取任务”动作串行,但任务处理仍可并行,适合演示固定 worker 模型。

7.1 数据模型

rust
use std::fmt;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{mpsc, Arc, Mutex};
use std::thread::{self, JoinHandle};

#[derive(Debug)]
struct ResizeJob {
    id: u64,
    source: String,
    width: u32,
}

#[derive(Debug)]
enum WorkerCommand {
    Process(ResizeJob),
    Stop,
}

#[derive(Debug)]
struct ResizeResult {
    id: u64,
    outcome: Result<String, ResizeError>,
}

#[derive(Debug)]
enum ResizeError {
    EmptySource,
    InvalidWidth,
}

impl fmt::Display for ResizeError {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        match self {
            Self::EmptySource => write!(f, "source path is empty"),
            Self::InvalidWidth => write!(f, "width must be greater than zero"),
        }
    }
}

impl std::error::Error for ResizeError {}

#[derive(Debug, Default, Clone)]
struct Summary {
    accepted: usize,
    succeeded: usize,
    failed: usize,
}

7.2 调度器结构

rust
struct ThumbnailPool {
    command_tx: mpsc::SyncSender<WorkerCommand>,
    result_rx: mpsc::Receiver<ResizeResult>,
    summary: Arc<Mutex<Summary>>,
    active: Arc<AtomicUsize>,
    workers: Vec<JoinHandle<()>>,
}

各字段只承担一个职责:

  • command_tx:有界入口,限制积压。
  • result_rx:把结果所有权交还调用者。
  • summary:维护多字段一致性。
  • active:提供近实时的独立执行数量。
  • workers:关闭时等待线程完成。

7.3 创建 worker

rust
impl ThumbnailPool {
    fn new(worker_count: usize, queue_capacity: usize) -> Self {
        assert!(worker_count > 0);

        let (command_tx, command_rx) =
            mpsc::sync_channel::<WorkerCommand>(queue_capacity);
        let (result_tx, result_rx) = mpsc::channel::<ResizeResult>();

        let command_rx = Arc::new(Mutex::new(command_rx));
        let summary = Arc::new(Mutex::new(Summary::default()));
        let active = Arc::new(AtomicUsize::new(0));
        let mut workers = Vec::with_capacity(worker_count);

        for worker_id in 0..worker_count {
            let command_rx = Arc::clone(&command_rx);
            let result_tx = result_tx.clone();
            let summary = Arc::clone(&summary);
            let active = Arc::clone(&active);

            workers.push(thread::spawn(move || loop {
                let command = {
                    let receiver = command_rx.lock().unwrap();
                    receiver.recv()
                };

                match command {
                    Ok(WorkerCommand::Process(job)) => {
                        active.fetch_add(1, Ordering::Relaxed);

                        let result = resize(job);

                        {
                            let mut snapshot = summary.lock().unwrap();
                            if result.outcome.is_ok() {
                                snapshot.succeeded += 1;
                            } else {
                                snapshot.failed += 1;
                            }
                        }

                        active.fetch_sub(1, Ordering::Relaxed);

                        if result_tx.send(result).is_err() {
                            eprintln!("worker {worker_id}: result receiver closed");
                            break;
                        }
                    }
                    Ok(WorkerCommand::Stop) | Err(_) => break,
                }
            }));
        }

        drop(result_tx);

        Self {
            command_tx,
            result_rx,
            summary,
            active,
            workers,
        }
    }

接收锁只在 recv() 期间持有;取得作业后立即释放,因此耗时处理不会占住该锁。需要更高吞吐时可以使用真正的 MPMC channel,但数据流和退出协议仍然相同。

7.4 提交、查询与关闭

rust
    fn submit(
        &self,
        job: ResizeJob,
    ) -> Result<(), mpsc::SendError<WorkerCommand>> {
        self.command_tx.send(WorkerCommand::Process(job))?;
        self.summary.lock().unwrap().accepted += 1;
        Ok(())
    }

    fn recv(&self) -> Result<ResizeResult, mpsc::RecvError> {
        self.result_rx.recv()
    }

    fn active(&self) -> usize {
        self.active.load(Ordering::Relaxed)
    }

    fn summary(&self) -> Summary {
        self.summary.lock().unwrap().clone()
    }

    fn shutdown(mut self) {
        for _ in &self.workers {
            let _ = self.command_tx.send(WorkerCommand::Stop);
        }

        drop(self.command_tx);

        for worker in self.workers.drain(..) {
            if worker.join().is_err() {
                eprintln!("worker panicked during shutdown");
            }
        }
    }
}

关闭协议向每个 worker 发送一个停止消息,再回收全部句柄。shutdown(self) 消费调度器,避免调用者在关闭后继续提交。

submit 中“发送成功后再增加 accepted”是一个可讨论的语义选择:它表示只有确实进入队列的任务才被接受。若锁 poisoned 或进程在两步之间失败,仍可能出现统计与事实不完全一致;要求事务级一致性时,应重新设计状态归属,而不是继续叠加锁。

7.5 业务处理函数

rust
fn resize(job: ResizeJob) -> ResizeResult {
    let outcome = if job.source.trim().is_empty() {
        Err(ResizeError::EmptySource)
    } else if job.width == 0 {
        Err(ResizeError::InvalidWidth)
    } else {
        Ok(format!("{}@{}px", job.source, job.width))
    };

    ResizeResult {
        id: job.id,
        outcome,
    }
}

7.6 使用示例

rust
fn main() {
    let pool = ThumbnailPool::new(3, 8);

    for (id, source, width) in [
        (1, "cover.jpg", 320),
        (2, "", 200),
        (3, "poster.png", 0),
        (4, "avatar.webp", 96),
    ] {
        pool.submit(ResizeJob {
            id,
            source: source.to_string(),
            width,
        })
        .unwrap();
    }

    for _ in 0..4 {
        let result = pool.recv().unwrap();
        println!("job {} => {:?}", result.id, result.outcome);
    }

    println!("active = {}", pool.active());
    println!("summary = {:?}", pool.summary());
    pool.shutdown();
}

结果抵达顺序不保证与提交顺序一致。需要稳定展示顺序时,应保留 id 并在汇总端排序,而不是给 worker 之间增加无必要的串行约束。


八、活性风险:正确访问内存还不够

8.1 死锁

两个线程按相反顺序获取两把锁,是最常见的死锁来源。应统一锁顺序,尽量避免同时持有多把锁,并禁止在锁内调用不受控制的回调。

8.2 饥饿与长临界区

即使没有死锁,一个线程长期持锁也可能让其他线程迟迟无法前进。缩短 guard 作用域、拆分状态和减少锁竞争,通常比盲目换成 RwLock 更有效。

8.3 忙等

rust
while !ready.load(Ordering::Acquire) {
    std::hint::spin_loop();
}

自旋只适合预期极短的等待。时间不可控时应使用 channel、Condvar 或其他阻塞原语,避免占满 CPU。

8.4 无界资源

无界 channel、每个请求创建一个线程、无限保存 JoinHandle 都会把流量峰值转化为内存或线程数量增长。并发系统必须同时设计容量、拒绝策略、超时和监控。


九、错误边界与全局初始化

线程基础设施错误与业务错误应分开表达:

  • channel 断开;
  • worker panic
  • 锁 poisoned;
  • 作业输入无效;
  • 作业执行失败。

业务错误随 JobResult 返回,基础设施错误由调度器 API 处理。这样调用者才能决定重试、降级还是终止。

运行期只初始化一次的共享配置可以使用 OnceLock

rust
use std::sync::OnceLock;

static WORKER_LIMIT: OnceLock<usize> = OnceLock::new();

fn worker_limit() -> usize {
    *WORKER_LIMIT.get_or_init(|| 4)
}

这比 static mut 更能表达一次初始化和跨线程可见性。


十、常见误区

  1. 线程创建后会自动妥善结束。 进程退出会终止未完成线程;关键线程必须管理句柄和退出协议。
  2. move 会复制捕获数据。 它只改变捕获所有权方式。
  3. channel 永远不会阻塞。 recv 会等待;有界队列满时 send 也会等待。
  4. Arc<T> 让内部类型自动线程安全。 Arc 只保护引用计数。
  5. 锁能消除所有竞态。 分成两次加锁的“检查后修改”仍可能产生逻辑竞态。
  6. Atomic 总比 Mutex 快。 复杂状态使用 Atomic 可能增加重试、缓存争用和证明成本。
  7. Send / Sync 表示程序不会死锁。 它们只描述内存安全层面的线程许可。
  8. 只要功能测试通过,并发代码就可靠。 时序错误往往需要压力测试、故障注入和退出路径测试才能暴露。

十一、实践练习

  1. 为缩略图调度器增加“拒绝新任务”和“等待队列清空”两种关闭模式。
  2. 把结果通道改为有界通道,观察结果消费过慢如何反向限制 worker。
  3. 给作业增加超时状态,并说明标准线程无法被安全强杀时应如何协作取消。
  4. Barrier 让多个线程同时开始,构造可重复的锁竞争测试。
  5. 保留一个 Sender clone,解释接收循环为何无法结束,再修复其所有权。
  6. Summary 的字段拆成 Atomic,编写一个读快照的线程,说明为什么可能观察到不一致组合。
  7. 使用 thread::scope 并行计算一个切片的多个不重叠分区,避免不必要的 Arc
  8. 为调度器定义 PoolError,区分提交失败、结果通道关闭和 worker panic
  9. 替换为 MPMC channel,并比较吞吐、依赖成本和关闭语义。

十二、总结

  1. 并发设计首先确定数据归属、同步关系和资源上限,而不是先选择线程 API。
  2. spawn 创建执行单元,move 转移捕获所有权,join 回收结果和故障。
  3. channel 适合移交任务与结果,有界通道还能表达背压。
  4. Arc 解决共享寿命,Mutex / RwLock 解决同步访问,两者职责不同。
  5. 锁应保护最小且完整的不变量,guard 不应跨越慢操作或未知回调。
  6. Atomic 适合独立、窄小的状态;复合一致性通常更适合锁。
  7. SendSync 是类型层面的线程安全许可,不保证活性、公平或业务一致性。
  8. 线程必须有明确退出协议,通道、句柄和所有 clone 的寿命都属于协议的一部分。
  9. 并发可靠性还需要处理过载、死锁、饥饿、poisoning、任务失败和停机路径。

十三、章节导航

上一章讨论单线程中的共享所有权与内部可变性;本章将数据流扩展到 OS 线程。下一章进入异步模型,观察等待如何被编译成 Future 状态机,并由运行时调度大量轻量任务。

  • 上一章:第8章:智能指针
  • 下一章:第10章:异步编程 —— Future、Pin、Stream 与任务调度

  • Rust 深度教程
  • 第7章:函数式 Rust
  • 第8章:智能指针
  • 第10章:异步编程
  • Claude Code Harness 深度教程

本章目录
一、先确定数据怎样流动二、线程的生命周期:spawn、move、join三、消息传递:让所有权沿协议移动四、共享状态:Arc 负责寿命,Mutex 负责访问五、Atomic:只在状态足够窄时使用六、Send 与 Sync:类型层面的跨线程许可七、实战:固定 worker 的缩略图调度器八、活性风险:正确访问内存还不够九、错误边界与全局初始化十、常见误区十一、实践练习十二、总结十三、章节导航Related Documents
苏ICP备2025204887号-2