Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

rust并发编程基础

1. 并发编程相关基础概念

1.1 并发和并行

  • 并发(Concurrency):指多个任务在同一时间段内交替执行(宏观同时,微观串行)。例如单核CPU上快速切换多个线程,给人一种“同时运行”的错觉。
  • 并行(Parallelism):指多个任务在同一时刻真正同时执行,需要多核CPU的支持。

💡 简单记忆:并发是“逻辑上的同时”,并行是“物理上的同时”。并发是并行的基础,但并发不一定并行。

1.2 进程和线程

  • 进程:操作系统资源分配的基本单位,拥有独立的内存空间、文件句柄等。进程间通信(IPC)开销较大。
  • 线程:CPU调度的基本单位,共享所属进程的内存和资源。线程间通信更高效,但也带来了数据竞争的风险。

💡 进程就像一个独立的房子,有自己的水电煤气。线程就像房子里不同房间的人,共用同一个房子的资源。

1.3 系统线程和绿色线程

  • 系统线程:由操作系统内核管理,创建、切换、销毁都需要内核参与,开销相对较大。Rust标准库中的 std::thread就是系统线程。
  • 绿色线程:由用户态运行时管理(如Go的goroutine),创建和切换开销极小,但需要运行时支持。Rust早期有绿色线程,但最终选择暴露系统线程作为标准。

Rust线程是什么:Rust标准库直接使用1:1模型,即一个Rust线程对应一个操作系统线程。这样做的优点是:与系统交互简单、无额外运行时开销、可预测的性能。对于需要大量轻量级任务,可以选择 tokioasync-std等异步运行时(基于绿色线程思想)。

2. 线程基本操作

2.1 创建新线程 – std::thread::spawn

功能:创建一个新的OS线程,并立即开始执行传入的闭包。

接口签名

#![allow(unused)]
fn main() {
pub fn spawn<F, T>(f: F) -> JoinHandle<T> 
where
    F: FnOnce() -> T + Send + 'static,
    T: Send + 'static,
}
  • 参数f – 一个闭包,在新线程中执行,闭包必须实现 Send + 'static(通常使用 move转移所有权)。
  • 返回值JoinHandle<T> – 代表新线程的句柄,可用于等待线程结束并获取返回值。

简单调用示例

#![allow(unused)]
fn main() {
let handle = std::thread::spawn(|| {
    println!("Hello from a new thread!");
});
}

2.2 等待线程结束 – JoinHandle::join

功能:阻塞当前线程,直到对应的线程执行完毕,并返回线程闭包的返回值。

接口签名

#![allow(unused)]
fn main() {
pub fn join(self) -> Result<T, Box<dyn Any + Send + 'static>>
}
  • 参数:无(消耗 self)。
  • 返回值Result<T, Box<dyn Any + Send + 'static>> – 成功返回 Ok(T),T是闭包返回值;失败返回 Err(Box<...>)包含panic信息。

简单调用示例

#![allow(unused)]
fn main() {
let handle = std::thread::spawn(|| 42);
let result = handle.join().unwrap();   // result = 42
}

2.3 线程闭包中的 move语义

功能:将闭包捕获的外部变量所有权转移到闭包内部,从而安全地在另一个线程中使用这些变量。

说明:不使用 move时,闭包会借用外部变量,但新线程可能存活超过变量所在作用域,导致悬垂引用。move强制转移所有权,确保变量在线程执行期间始终有效。

简单调用示例

#![allow(unused)]
fn main() {
let data = vec![1, 2, 3];
let handle = std::thread::spawn(move || {
    println!("{:?}", data);   // data所有权已移入线程
});
// println!("{:?}", data); // 编译错误:data已移动
handle.join().unwrap();
}

3. 线程同步

3.1 什么是线程同步

当多个线程同时访问共享数据时,为了防止数据竞争、保证一致性和正确性,需要采用某种协调机制,这就是线程同步。例如:互斥锁、条件变量、信号量、读写锁等。

3.2 为什么需要线程同步 – 无同步的例子

考虑一个没有同步的例子:两个线程同时对一个共享计数器进行递增操作。

#![allow(unused)]
fn main() {
static mut COUNTER: i32 = 0;

let t1 = std::thread::spawn(|| {
    for _ in 0..1000 {
        unsafe { COUNTER += 1; }
    }
});
let t2 = std::thread::spawn(|| {
    for _ in 0..1000 {
        unsafe { COUNTER += 1; }
    }
});
t1.join().unwrap();
t2.join().unwrap();
unsafe { println!("{}", COUNTER); } // 可能不是2000,而是例如1987等随机值
}

由于线程交错执行,COUNTER += 1并非原子操作(实际上是读-改-写三步),导致丢失更新,结果不可预测。这就是典型的数据竞争,需要同步来避免。

3.3 线程同步常用手段

  • 互斥锁(Mutex):保证同一时刻只有一个线程访问数据。
  • 读写锁(RwLock):允许多个读线程或一个写线程同时访问。
  • 条件变量(Condvar):让线程等待某个条件满足后再继续执行。
  • 原子类型(Atomic):对简单类型提供无锁的原子操作。
  • 消息传递(Channel):通过发送/接收消息进行线程间通信,Rust中的 mpsc是典型。

3.4 Rust互斥锁(Mutex<T>

3.4.1 Rust的 Mutex与Java/C++的不同

  • Java/C++:互斥锁通常是一个独立对象,你需要手动 lock()/unlock(),并且锁与数据是分离的。容易忘记释放锁,或者持有锁时间过长。
  • RustMutex<T>是一个容器,它包裹了数据 T。你无法直接访问内部数据,而必须通过 lock()方法得到一个 MutexGuard(类似智能指针),该守卫在作用域结束时自动释放锁。这种设计保证了数据受锁的保护,且不会忘记解锁。

核心思想:Rust的 Mutex<T>RefCell<T>类似,都提供内部可变性——即通过不可变引用也能修改内部数据。只是 RefCell在运行时检查借用规则,而 Mutex通过阻塞线程来保证独占访问。

3.4.2 创建互斥锁 – Mutex::new

功能:创建一个新的互斥锁,内部包裹初始值 data

接口签名

#![allow(unused)]
fn main() {
pub fn new(data: T) -> Mutex<T>
}
  • 参数data – 需要被保护的数据。
  • 返回值Mutex<T>实例。

简单调用示例

#![allow(unused)]
fn main() {
use std::sync::Mutex;
let m = Mutex::new(100);
}

3.4.3 加锁 – Mutex::lock

功能:阻塞当前线程,直到获得互斥锁,返回一个 MutexGuard智能指针。如果持有锁的线程panic,lock会返回错误。

接口签名

#![allow(unused)]
fn main() {
pub fn lock(&self) -> LockResult<MutexGuard<'_, T>>
}
  • 参数:无(通过 &self获取锁)。
  • 返回值LockResult<MutexGuard<'_, T>>,通常调用 .unwrap()获取 MutexGuard。该守卫实现了 DerefDerefMut,可直接访问内部数据;离开作用域时自动解锁。

简单调用示例

#![allow(unused)]
fn main() {
let guard = m.lock().unwrap();   // 获取锁
*guard += 1;                     // 修改内部值
println!("{}", *guard);          // 读取
// guard 离开作用域,自动解锁
}

3.4.4 多线程共享 Mutex<T> – 配合 Arc

由于 Mutex<T>本身不实现 Copy,且多个线程需要共享所有权,因此通常将 Mutex<T>放入 Arc(原子引用计数)中。Arc允许多个线程同时拥有同一个 Mutex的所有权。

简单调用示例

#![allow(unused)]
fn main() {
use std::sync::{Arc, Mutex};
use std::thread;

let counter = Arc::new(Mutex::new(0));
let mut handles = vec![];

for _ in 0..10 {
    let counter = Arc::clone(&counter);
    let handle = thread::spawn(move || {
        let mut num = counter.lock().unwrap();
        *num += 1;
    });
    handles.push(handle);
}
for handle in handles {
    handle.join().unwrap();
}
println!("Result: {}", *counter.lock().unwrap());
}

3.5 条件变量(Condvar

条件变量用于线程间的等待-通知机制:一个线程等待某个条件成立,另一个线程满足条件后发出通知。Rust的 Condvar总是与 Mutex配合使用。

3.5.1 创建条件变量 – Condvar::new

功能:创建一个新的条件变量。

接口签名

#![allow(unused)]
fn main() {
pub fn new() -> Condvar
}
  • 参数:无。
  • 返回值Condvar实例。

简单调用示例

#![allow(unused)]
fn main() {
use std::sync::{Arc, Mutex, Condvar};
let pair = Arc::new((Mutex::new(false), Condvar::new()));
}

3.5.2 条件等待 – Condvar::wait_while

功能:在持有 Mutex锁的情况下,检查等待条件(闭包)。如果闭包返回 true,表示“还需要继续等待”,它会自动释放锁并阻塞当前线程;当被其他线程通知后重新获取锁并再次检查。循环直到闭包返回 false

接口签名

#![allow(unused)]
fn main() {
pub fn wait_while<'a, T, F>(
    &self,
    guard: MutexGuard<'a, T>,
    condition: F
) -> LockResult<MutexGuard<'a, T>>
where
    F: FnMut(&mut T) -> bool,
}
  • 参数
    • guard – 已经获得的锁守卫。
    • condition – 一个闭包,接收 &mut T,返回 bool。返回 true 表示继续等待,返回 false 表示停止等待并返回锁守卫。
  • 返回值:返回一个新的 MutexGuard<'a, T>,此时等待条件已经不再成立。

简单调用示例

#![allow(unused)]
fn main() {
let mut guard = lock.lock().unwrap();
guard = condvar.wait_while(guard, |data| data.is_empty()).unwrap();
// 现在 guard 中的队列非空
}

3.5.3 唤醒等待线程 – Condvar::notify_one / Condvar::notify_all

功能:唤醒一个或所有等待在该条件变量上的线程。

接口签名

#![allow(unused)]
fn main() {
pub fn notify_one(&self)
pub fn notify_all(&self)
}
  • 参数:无。
  • 返回值:无。

简单调用示例

#![allow(unused)]
fn main() {
condvar.notify_one();   // 唤醒一个线程
condvar.notify_all();   // 唤醒所有等待线程
}

3.6 多生产者单消费者通道(mpsc

mpsc代表“Multiple Producer, Single Consumer”。它是Rust标准库提供的一个消息传递同步工具,允许多个线程发送消息,但只有一个线程接收消息。

3.6.1 创建通道 – std::sync::mpsc::channel

功能:创建一个新的异步通道,返回发送端和接收端。发送端可以克隆(多生产者),接收端独占。

接口签名

#![allow(unused)]
fn main() {
pub fn channel<T>() -> (Sender<T>, Receiver<T>)
}
  • 参数:无(通过泛型 T指定消息类型)。
  • 返回值:元组 (Sender<T>, Receiver<T>)

简单调用示例

#![allow(unused)]
fn main() {
use std::sync::mpsc;
let (tx, rx) = mpsc::channel();
}

3.6.2 发送消息 – Sender::send

功能:将消息发送到通道。如果接收端已经关闭,返回错误。

接口签名

#![allow(unused)]
fn main() {
pub fn send(&self, t: T) -> Result<(), SendError<T>>
}
  • 参数t – 要发送的消息(消耗所有权)。
  • 返回值:成功返回 Ok(()),失败返回 Err(SendError(t))(将消息返回)。

简单调用示例

#![allow(unused)]
fn main() {
tx.send(42).unwrap();
}

3.6.3 接收消息 – Receiver::recv

功能:阻塞当前线程,直到通道中有消息可接收,或所有发送端已关闭(此时返回错误)。

接口签名

#![allow(unused)]
fn main() {
pub fn recv(&self) -> Result<T, RecvError>
}
  • 参数:无。
  • 返回值:成功返回 Ok(T),失败返回 Err(RecvError)(表示没有更多消息)。

简单调用示例

#![allow(unused)]
fn main() {
let msg = rx.recv().unwrap();
println!("Received: {}", msg);
}

3.6.4 完整示例片段

#![allow(unused)]
fn main() {
use std::sync::mpsc;
use std::thread;

let (tx, rx) = mpsc::channel();
thread::spawn(move || {
    tx.send("Hello from thread").unwrap();
});
println!("{}", rx.recv().unwrap());
}

4. 综合示例代码:生产者消费者模型实现

下面我们使用前面介绍的 MutexCondvarArcAtomicBool实现一个安全、可停止的生产者-消费者队列。该示例完整展示了线程同步的多项技术。

4.1 代码结构分析

  • Queue<T>:队列数据结构,内部包含:
    • data: Mutex<VecDeque<T>> – 互斥锁保护的队列。
    • condvar: Condvar – 条件变量,用于等待队列非空或停止信号。
    • stopped: AtomicBool – 原子标志位,指示是否停止生产/消费。
  • push:生产者调用,添加数据并通知消费者。
  • pop:消费者调用,如果队列为空则等待,直到有数据或停止信号。
  • stop:设置停止标志并唤醒所有等待线程。

4.2 完整代码

#![allow(unused)]
fn main() {
use std::collections::VecDeque;
use std::sync::{Arc, Mutex, Condvar, atomic::{AtomicBool, Ordering}};

struct Queue<T> {
    data: Mutex<VecDeque<T>>,
    condvar: Condvar,
    stopped: AtomicBool,
}

impl<T> Queue<T> {
    fn new() -> Arc<Queue<T>> {
        Arc::new(Queue {
            data: Mutex::new(VecDeque::new()),
            condvar: Condvar::new(),
            stopped: AtomicBool::new(false),
        })
    }

    fn stop(&self) {
        self.stopped.store(true, Ordering::Release);
        self.condvar.notify_all();
    }

    fn push(&self, data: T) {
        let mut guard = self.data.lock().unwrap();
        if self.stopped.load(Ordering::Acquire) { return; }
        guard.push_back(data);
        self.condvar.notify_one();
    }

    fn pop(&self) -> Option<T> {
        let mut guard = self.data.lock().unwrap();
        // 等待条件:队列非空 或者 已经停止
        guard = self.condvar.wait_while(guard, |g| {
            g.is_empty() && !self.stopped.load(Ordering::Acquire)
        }).unwrap();

        if guard.is_empty() && self.stopped.load(Ordering::Acquire) {
            None
        } else {
            guard.pop_front()
        }
    }
}

pub fn demo() {
    println!("...............生产者消费者模型示例开始...................");
    let queue = Queue::<i32>::new();
    let produce_queue = queue.clone();
  
    let produce = std::thread::spawn(move || {
        let mut count = 0;
        while count < 11 {
            produce_queue.push(count);
            count += 1;
        }
        produce_queue.stop();
    });

    let consume_queue1 = queue.clone();
    let consume1 = std::thread::spawn(move || {
        while let Some(data) = consume_queue1.pop() {
            println!("consumer1 consume data: {}", data);
        }
    });

    let consume_queue2 = queue.clone();
    let consume2 = std::thread::spawn(move || {
        while let Some(data) = consume_queue2.pop() {
            println!("consumer2 consume data: {}", data);
        }
    });

    produce.join().unwrap();
    consume1.join().unwrap();
    consume2.join().unwrap();
    println!("...............生产者消费者模型示例结束...................");
}
}