首页
/ Rustlings 线程练习详解:thread::spawn、Arc<Mutex<T>> 与 mpsc 消息传递的三步进阶

Rustlings 线程练习详解:thread::spawn、Arc<Mutex<T>> 与 mpsc 消息传递的三步进阶

2026-09-04 19:50:42作者:廉皓灿Ida

本篇技术指南以 Rustlings 的 Threads 练习组 为主体,系统讲解 Rust 中多线程编程的三个核心机制:用 thread::spawn 创建线程并通过 JoinHandle 回收返回值、用 Arc<Mutex<T>> 安全共享可变状态、用 std::sync::mpsc 在线程间传递消息。读完并动手完成 threads1threads3 三道练习后,你将具备在 Rust 中编写正确、可编译、无数据竞争并发代码的基本能力。

1. 线程与进程:练习组的概念起点

Threads 练习组的 README 首先给出了并发模型的基础定义:

  • 在当前大多数操作系统中,一个程序被打包为进程(process),代码在进程内运行,操作系统同时管理多个进程;
  • 进程内部还可以存在同时独立运行的部分,承载这些独立执行单元的特性就叫线程(thread)

这一定义对应着 Rust 官方指南(The Rust Book)第 16 章的主题脉络:从“用线程同时运行代码”到“用消息传递在线程间转移数据”。Rustlings 的三道练习恰好沿这条脉络递进:threads1 让你接触线程的生命周期管理,threads2 引入共享状态的同步原语,threads3 则切换到 Rust 推崇的“转移所有权代替共享”的消息传递模型。

练习文件均位于 exercises/20_threads/ 目录,官方参考答案位于 solutions/20_threads/ 目录,可对照学习。

2. threads1:用 JoinHandle 等待线程并收集返回值

threads1.rs 的任务是:启动 10 个线程,每个线程至少运行 250ms,并以“实际耗时(毫秒数)”作为返回值;主线程必须等待所有线程结束,并把 10 个返回值收集进一个 Vec

练习给出的完整框架代码如下(thread::spawn、睡眠计时、打印逻辑均已写好):

use std::{
    thread,
    time::{Duration, Instant},
};

fn main() {
    let mut handles = Vec::new();
    for i in 0..10 {
        let handle = thread::spawn(move || {
            let start = Instant::now();
            thread::sleep(Duration::from_millis(250));
            println!("Thread {i} done");
            start.elapsed().as_millis()
        });
        handles.push(handle);
    }

    let mut results = Vec::new();
    for handle in handles {
        // TODO: Collect the results of all threads into the `results` vector.
        // Use the `JoinHandle` struct which is returned by `thread::spawn`.
    }

    if results.len() != 10 {
        panic!("Oh no! Some thread isn't done yet!");
    }

    println!();
    for (i, result) in results.into_iter().enumerate() {
        println!("Thread {i} took {result}ms");
    }
}

关键点在于 thread::spawn 的返回值:它返回一个 JoinHandle<T>,其中 T 是闭包的返回类型。JoinHandle 提供两个能力——join() 会阻塞等待该线程结束并返回 Result<T, Box<dyn Any + Send>>(线程正常结束则为 Ok(返回值),线程内部发生 panic 则为 Err)。

对照 参考解答,TODO 处只需一行:

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

这里有两点值得注意:

  1. 必须先 join 再读结果。如果不等待线程完成就检查 results.len() != 10,主线程可能提前 panic——join() 正是“等待所有已创建线程结束”这一需求的直接实现;
  2. unwrap() 的语义join() 失败只在子线程 panic 时发生。本练习中子线程只做睡眠与打印,不会 panic,因此 unwrap() 是安全写法;在真实项目中,通常应对 Err 做日志记录或错误传播。

3. threads2:Arc 不够,共享可变状态需要 Arc<Mutex>

threads2.rs 在上一步基础上进一步要求:10 个线程要共同更新一个共享值 JobStatus.jobs_done,全部完成后打印其最终值。

练习的初始代码是这样的:

use std::{sync::Arc, thread, time::Duration};

struct JobStatus {
    jobs_done: u32,
}

fn main() {
    // TODO: `Arc` isn't enough if you want a **mutable** shared state.
    let status = Arc::new(JobStatus { jobs_done: 0 });

    let mut handles = Vec::new();
    for _ in 0..10 {
        let status_shared = Arc::clone(&status);
        let handle = thread::spawn(move || {
            thread::sleep(Duration::from_millis(250));

            // TODO: You must take an action before you update a shared value.
            status_shared.jobs_done += 1;
        });
        handles.push(handle);
    }

    // Waiting for all jobs to complete.
    for handle in handles {
        handle.join().unwrap();
    }

    // TODO: Print the value of `JobStatus.jobs_done`.
    println!("Jobs done: {}", todo!());
}

这个练习直接暴露了一个新手常见误区:Arc(原子引用计数)只解决了“多线程共享所有权”的问题,它本身不保证共享的引用是可变访问。status_shared.jobs_done += 1 要求通过共享引用进行写操作,Rust 的所有权规则在编译期就会拒绝它。

solutions/20_threads/threads2.rs 给出了标准答案,核心改动有三处:

// 1. 用 Mutex 包裹共享状态
let status = Arc::new(Mutex::new(JobStatus { jobs_done: 0 }));

// 2. 线程内先加锁,再更新
let status_shared = Arc::clone(&status);
let handle = thread::spawn(move || {
    thread::sleep(Duration::from_millis(250));
    status_shared.lock().unwrap().jobs_done += 1;
});

// 3. 主线程读值时同样先加锁
println!("Jobs done: {}", status.lock().unwrap().jobs_done);

从源码结构看,这个组合的含义是:

  • Mutex::new(JobStatus { .. }) 让“数据 + 互斥访问”成为一体,任何时刻只有一方能持有锁并修改 jobs_done,从而杜绝数据竞争;
  • status_shared.lock() 返回 Result<MutexGuard, PoisonError>unwrap() 在锁被“投毒”(持锁线程 panic)时会传播该错误。练习中线程无 panic 路径,因此 unwrap() 可接受;
  • Arc::clone(&status) 增加引用计数,把同一份堆上数据的访问权移交给每个闭包(move || 闭包)。

最终 jobs_done 的值应为 10——每个线程恰好加一。

4. threads3:mpsc 消息传递与 Sender 的克隆

threads3.rs 切换到 Rust 的另一条并发路线——消息传递(message passing):不共享数据,而是把数据的所有权沿着 channel 转移出去。

题目结构如下:一个 Queue 持有两半数据(1~5 与 6~10),send_tx 函数把两半分别交给两个线程发送:

fn send_tx(q: Queue, tx: mpsc::Sender<u32>) {
    // TODO: We want to send `tx` to both threads. But currently, it is moved
    // into the first thread. How could you solve this problem?
    thread::spawn(move || {
        for val in q.first_half {
            println!("Sending {val:?}");
            tx.send(val).unwrap();
            thread::sleep(Duration::from_millis(250));
        }
    });

    thread::spawn(move || {
        for val in q.second_half {
            println!("Sending {val:?}");
            tx.send(val).unwrap();
            thread::sleep(Duration::from_millis(250));
        }
    });
}

问题出在 tx 只能被移动进一个闭包:第二个 thread::spawn 中的 tx.send(val) 会因 tx 已被第一个闭包移走而编译失败。这与 threads2Arc 的困境形成对照——同样是“一份资源要分给多个使用者”,共享路线靠 Arc 复制引用计数,消息传递路线则靠 Sender 自身的 clone。参考解答只加了一行:

fn send_tx(q: Queue, tx: mpsc::Sender<u32>) {
    let tx_clone = tx.clone();
    thread::spawn(move || {
        for val in q.first_half {
            println!("Sending {val:?}");
            tx_clone.send(val).unwrap();
            thread::sleep(Duration::from_millis(250));
        }
    });
    // 第二个线程使用原始 `tx`
    thread::spawn(move || {
        for val in q.second_half {
            println!("Sending {val:?}");
            tx.send(val).unwrap();
            thread::sleep(Duration::from_millis(250));
        }
    });
}

mpsc(multi-producer, single-consumer)channel 的设计天然支持多发送端:Sender<T> 可自由 clone,每个克隆都是独立的生产者;而 Receiver<T> 只能有一个。所有 Sender(含克隆)被丢弃后,Receiver 会感知到通道关闭,迭代结束——这正是练习自带测试的验证方式。

值得一提的是,本练习的验证不走 main,而是文件末尾的单元测试 tests 模块

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn threads3() {
        let (tx, rx) = mpsc::channel();
        let queue = Queue::new();

        send_tx(queue, tx);

        let mut received = Vec::with_capacity(10);
        for value in rx {
            received.push(value);
        }

        received.sort();
        assert_eq!(received, [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]);
    }
}

测试中 for value in rx 依赖“所有 Sender 已随 send_tx 调用结束而转移进线程”这一事实:一旦两个线程都结束,通道关闭,rx 迭代器返回 None,主测试线程得以继续执行断言。由于两个线程的发送顺序不确定,测试先对 received 排序再与 1..=10 比较,规避了乱序带来的误报。

5. 三种机制的选择与后续学习

三道练习覆盖了 Rust 并发的两条基本路线:

机制 代表练习 适用场景
thread::spawn + JoinHandle::join threads1 需要等待并发任务完成并回收其返回值
Arc<Mutex<T>> 共享可变状态 threads2 多个线程需要读写同一份数据
mpsc channel 消息传递 threads3 生产者/消费者结构,数据所有权单向流动

练习组 README 在“Further information”中指向了 Rust 官方文档的三个延伸阅读方向,可作为本练习组的理论补充:

  • The Rust Book 第 16.1 节“Using Threads to Run Code Simultaneously”——对应 threads1spawn/join 模型;
  • The Rust Book 第 16.2 节“Using Message Passing to Transfer Data Between Threads”——对应 threads3mpsc 模型;
  • “Dining Philosophers”示例——经典的互斥锁与死锁问题演示,理解 threads2Mutex 语义后值得深入。

完成本练习组后,可继续前往 19_smart_pointersRc/RefCell 单线程共享)或 17_teststhreads3 所用的测试模块写法)等相邻章节巩固相关知识。所有练习均可通过 rustlings CLI 在本仓库中逐题验证,参考解答统一存放于 solutions/ 目录,便于自查对比。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
527
590
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
904
1.82 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
889
5.78 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.52 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.33 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
981
502
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384