Rust: Rust 并发编程:线程、消息传递与共享状态

最后更新:2026-08-26

并发(Concurrency)是程序"同时"做多件事的能力——Rust 通过所有权系统和类型系统在编译时消除数据竞争,让你写出既高效又安全的并发代码。

如果说单线程程序是"一个人在厨房里从头做到尾",那多线程就是"多个厨师同时工作"——有人切菜、有人炒菜、有人装盘。但厨房里人多了就容易出乱子:两个人同时用一把刀、一个人拿了另一个人正在用的配料。Rust 的并发模型就是"带规则的多厨师厨房"——每个厨师有自己的工具,传递食材走专用通道,共享的调味料只能一个人用。


1. 你将学到


2. 多人协作文档编辑系统的故事

(1) 痛苦:没有并发控制的在线文档

Luna 所在的公司正在开发一个在线的协作文档编辑器。早期版本只允许一个人编辑,团队成员抱怨不断:

"如果能给每个编辑者开一个独立线程、修改请求通过通道传递、文档内容用互斥锁保护就好了……"

(2) Rust 并发模型的方案

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 保护。ArcMutex 可以在多个线程间共享。所有权系统确保不会出现数据竞争。


3. 核心概念

(1) Rust 并发编程体系

100%
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&lt;T&gt; 互斥锁"]
    D --> D2["lock() 获取锁"]
    D --> D3["Arc&lt;T&gt; 原子引用计数"]

    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 不靠运行时检查保证线程安全,而是通过 SendSync 这两个 marker trait 在编译时检查。如果一个类型不是 Send,你试图把它传到另一个线程就会编译报错。


4. 并发编程示例

▶ 示例 1:thread::spawn + join——创建和等待线程(难度 ⭐)

RUST
// ============================================
// 演示: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!("所有工作完成!");
}

输出:

TEXT 📖 仅展示
=== 厨房开工 ===
[厨师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——消息传递(难度 ⭐⭐)

RUST
// ============================================
// 演示: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)

输出:

TEXT 📖 仅展示
=== 协作文档编辑器 ===
[服务器] 系统 插入了: 自动保存文档
[服务器] 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>>——共享状态的线程安全访问(难度 ⭐⭐⭐)

RUST
// ============================================
// 演示: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);
}

输出:

TEXT 📖 仅展示
=== 共享计数器演示 ===
[工人0] 当前计数: 1
[工人0] 当前计数: 2
[工人1] 当前计数: 3
[工人1] 当前计数: 4
[工人0] 当前计数: 5
...(中间输出因线程调度而不同)
[工人4] 当前计数: 50
[工人4] 工作完成

=== 最终结果 ===
计数器最终值: 50
预期值: 50 (5个线程 × 10次)

Mutex<T> 提供互斥访问——同一时间只有一个线程能访问内部数据。lock() 返回 MutexGuard<T>,它实现了 DerefDrop:离开作用域时自动释放锁。Arc<T>(Atomic Reference Counted)是 Rc<T> 的多线程版本,使用原子操作保证引用计数的线程安全。Arc::clone 增加引用计数而不拷贝数据。


▶ 示例 4:线程间协作——生产者-消费者模式(难度 ⭐⭐⭐)

RUST
// ============================================
// 生产者-消费者:多生产者 + 单消费者
// ============================================

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!("所有生产者已退出,消费者结束");
}

输出(顺序可能不同):

TEXT 📖 仅展示
=== 消费者接收 ===
  收到: 生产者1: 苹果
  收到: 生产者2: 西瓜
  收到: 生产者1: 香蕉
  收到: 生产者1: 橙子
  收到: 生产者2: 葡萄
生产者1 完成
生产者2 完成
所有生产者已退出,消费者结束

tx.clone() 创建多个发送端,实现多生产者模式。当所有发送端(txtx2)被 drop 后,rx 迭代器自动结束——不需要手动发送"结束"信号。


▶ 示例 5:综合练习——并行数据处理(难度 ⭐⭐⭐)

RUST
// ============================================
// 并行计算:多线程分片求和
// ============================================

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);
}

输出:

TEXT 📖 仅展示
=== 串行求和 ===
1 到 1000 的和: 500500

=== 并行求和 (4 线程) ===
各线程部分和: [78126, 218874, 109374, 94126]
并行求和结果: 500500

=== 10M 数据性能对比 ===
串行: [时间]
并行 (8线程): [时间]

parallel_sum 将数据分片,每个线程计算部分和,通过 Arc<Mutex<i64>> 累加到总结果。chunk_size 向上取整确保所有数据被覆盖。assert_eq! 验证并行结果与串行一致。对于大数据集,多线程可显著提升性能。


❓ 常见问题

Q thread::spawn 创建的线程和主线程是什么关系?
A 所有线程是并行的,主线程结束整个程序就结束。
Q mpsc::channelsendrecv 会阻塞吗?
A send 通常不阻塞(缓冲区有空位时立即返回),recv 会阻塞直到收到消息。
Q Mutex<T>Arc<Mutex<T>> 必须一起用吗?
A 不一定。
Q 使用 Mutex 会不会有死锁?
A Rust 不阻止死锁,需要在代码中避免。
Q SendSync 有什么区别?
A Send 是"所有权可以跨线程转移",Sync 是"引用可以跨线程共享"。

📖 小节


📝 作业

  1. 难度 ⭐:写一个程序,创建 3 个线程分别计算 1..=10、11..=20、21..=30 的和,主线程用 join 等待所有线程结束后汇总三个部分和并打印最终结果。
  2. 难度 ⭐⭐:使用 mpsc::channel 实现一个"任务分发器"。创建 1 个生产者线程(生成 10 个任务编号 1-10)和 2 个消费者线程(每个消费者从通道接收任务编号并打印"[消费者X] 处理任务 #N")。确保所有任务都被处理完毕。
  3. 难度 ⭐⭐⭐:实现一个"共享银行账户"系统。使用 Arc<Mutex<f64>> 作为共享余额。创建 4 个线程分别模拟存款操作(每个线程随机存入 10-100 元)。主线程等待所有存款线程结束后读取最终余额。额外要求:在每次存款前后打印余额变化,验证没有数据竞争(最终余额 = 初始余额 + 所有存款之和)。
Web-Tutorial.com

Web-Tutorial 技术团队

由多位开发者共同维护的编程教程平台。每篇教程由对应领域的开发者编写和审核,确保内容准确可靠。如发现任何问题,欢迎向我们反馈。

100%

🙏 帮我们做得更好

我们是刚上线的编程教程站,几个人的小团队,精力有限。页面虽经检查,难免还有疏漏——链接失效、排版错乱、内容有误、语言生硬……

如果您发现了,麻烦告诉我们,我们会在收到反馈后第一时间进行修复,再次感谢您的光临 🙏