Rust 并发:任务、消息与共享状态

Concurrency in Rust

3,642 words 24 min read
目录 23 节
  1. 任务如何运行
  2. 先从线程开始
  3. move 移动数据
  4. 当线程只活在一个作用域里
  5. 编译器凭什么允许数据跨线程?
  6. 难道一个元素创建一条线程?
  7. 如果任务大部分时间都在等待
  8. 任务如何交换数据
  9. JoinHandle 只适合在汇合点拿结果
  10. 通道:把数据交给另一个任务
  11. 多个生产者,以及通道如何结束
  12. 无界队列把压力藏进了内存
  13. 通道可以把程序连成流水线
  14. 数据无法转移时:共享状态
  15. 有些数据确实需要共同访问
  16. Arc 解决共同所有权
  17. Mutex 解决互斥访问
  18. 读多写少就一定用 RwLock 吗?
  19. 一个计数器需要整把锁吗?
  20. 锁真正麻烦的地方
  21. 一个融合的例子
  22. 到底该从哪里开始选?
  23. 最后还是所有权

Rust 的并发工具有很多:标准库里有 threadchannelMutexRwLockAtomic;三方库还有 Rayon、Tokio。

一个并发程序真正需要回答的,通常是两件事:

  1. 工作如何被拆分和调度?
  2. 并发执行单元之间如何交换数据?

第二个问题还可以继续往下分:数据能不能直接交给另一个任务,还是必须由多个任务共同访问?

本文无意从 API 清单开始。我们从最普通的单个线程出发,看看问题是怎么出现的,又是如何解决的。

任务如何运行

先解决第一个问题:工作如何被拆分和调度。

先从线程开始

假设有两段彼此独立的计算:

fn calculate_left() -> u64 {
    // 一段耗时计算
    40
}

fn calculate_right() -> u64 {
    // 另一段耗时计算
    2
}

fn main() {
    let left = calculate_left();
    let right = calculate_right();

    println!("{}", left + right);
}

现在它们依次执行。既然两段计算互不依赖,自然会想到让其中一段去另一条线程运行:

use std::thread;

fn main() {
    let left = thread::spawn(calculate_left);
    let right = calculate_right();

    let result = left.join().unwrap() + right;
    println!("{result}");
}

thread::spawn 创建一条操作系统线程,并返回一个 JoinHandle。调用 join() 会等待那条线程结束,并拿回它的返回值。

如果把这段代码画出来,广义上它已经形成了一个最简单的 Fork-Join:主线程派生出一个并发分支,自己继续计算,随后在 join() 处重新汇合。

             ┌── calculate_left  ──┐
main ─ fork ─┤                     ├─ join ─ result
             └── calculate_right ──┘

Fork-Join 是一种任务结构:

  1. 把工作拆成几个可以独立执行的部分
  2. 让它们并发运行
  3. 等待子任务结束
  4. 汇总结果

这里的 spawn 是 fork,join 是 join。

join() 返回的是 Result。子线程正常结束时,里面是返回值;子线程发生 panic 时,错误里会带回 panic 的 payload。线程里的失败不会凭空消失,它会在聚合点重新出现。

为了突出并发结构,后面的示例会经常使用 unwrap()。实际程序通常需要区分线程 panic、通道断开、锁中毒和业务错误,再决定传播、恢复还是终止。

move 移动数据

上面的函数没有捕获外部数据。现实中的任务通常并非如此。

比如我们想在线程里读取一个 Vec

use std::thread;

fn main() {
    let values = vec![1, 2, 3, 4];

    let handle = thread::spawn(|| {
        println!("{values:?}");
    });

    handle.join().unwrap();
}

这段代码无法通过编译。

问题不在于 Vec 不能跨线程,而在于闭包默认借用了 values。但是子线程可能比 main 里的当前作用域活得更久。如果主线程提前结束并释放 values,子线程里的引用就会悬空。

虽然我们知道马上会调用 join(),但普通 thread::spawn 的签名并不依赖,也不应该依赖这种外部行为。它要求传入的闭包满足 'static,不能随便借用当前栈帧里的数据。

最直接的办法是加上 move

use std::thread;

fn main() {
    let values = vec![1, 2, 3, 4];

    let handle = thread::spawn(move || {
        println!("{values:?}");
    });

    handle.join().unwrap();
}

这次 values 的所有权被移动进闭包,再跟着闭包进入子线程。生命周期问题消失了,因为数据现在归子线程所有。

但代价也很明确:主线程不能再使用原来的 values

let handle = thread::spawn(move || {
    println!("{values:?}");
});

println!("{}", values.len()); // values 已经被移动

move 不是用来消除编译错误的咒语。它改变了数据的归属。

如果子线程就应该独立拥有这份数据,那么移动所有权很合理。但有些时候,我们只是想暂时借用当前函数里的数据,并且能够证明子线程不会离开这个函数。

这正是 thread::scope 要解决的问题。

当线程只活在一个作用域里

假设我们要并行处理一个切片的左右两半:

use std::thread;

fn double(values: &mut [i32]) {
    let mid = values.len() / 2;
    let (left, right) = values.split_at_mut(mid);

    thread::scope(|scope| {
        scope.spawn(|| {
            for value in left {
                *value *= 2;
            }
        });

        scope.spawn(|| {
            for value in right {
                *value *= 2;
            }
        });
    });
}

fn main() {
    let mut values = vec![1, 2, 3, 4];
    double(&mut values);

    assert_eq!(values, [2, 4, 6, 8]);
}

这里的两个闭包都借用了当前函数里的数据,而且还是可变借用。

它之所以安全,是因为 thread::scope 做出了一个普通 spawn 没有做出的保证:

scope 返回以前,里面启动的线程一定已经结束。

leftright 不会在线程仍然运行时被释放。另一方面,split_at_mut 已经把原切片拆成两个不重叠的可变切片,两个线程也不会同时修改同一块内存。

Rust 没有简单地说“跨线程不能借用”。它要求的是:你得给出一个足够明确的生命周期边界,让编译器能够证明这次借用不会逃出去。

普通 thread::spawn 适合把一份独立拥有的数据交给线程。thread::scope 则适合让几条生命周期明确的线程临时借用当前数据。

到这里,我们已经能处理固定数量的异构任务了。

编译器凭什么允许数据跨线程?

thread::scope 解决的是生命周期问题:子线程会不会活得比它借用的数据更久。

但还有另一个问题:这份数据本身能不能安全地跨过线程边界?

thread::spawn 的签名稍微简化一下,可以看到两个熟悉但容易被忽略的约束:

pub fn spawn<F, T>(f: F) -> JoinHandle<T>
where
    F: FnOnce() -> T + Send + 'static,
    T: Send + 'static,

'static => “这份数据能活多久”,Send => “它的所有权能不能安全地转移到另一条线程”。

Rust 还有一个经常与 Send 一起出现的 trait:Sync。如果 TSync,就意味着多个线程可以安全地共享 &T。换一种更准确的说法:

当且仅当 &TSend 时,T 才是 Sync

这两个 trait 通常不需要手动实现。编译器会根据一个类型的组成部分,推导它是否满足相应条件。

例如,Rc<T> 使用非原子的引用计数,不能安全地在线程之间转移;RefCell<T> 把借用检查放到了运行时,却没有提供线程间同步,因此也不是 Sync

Arc<T> 使用原子引用计数,可以解决跨线程的共同所有权,但它并不会让内部的 T 凭空变得线程安全。Arc<T> 能否跨线程共享,仍然取决于 T 自己是否满足相应的 SendSync 约束。需要修改共享数据时,通常还要借助 Mutex<T> 之类的同步原语来协调访问。

生命周期说明“能活多久”,Send / Sync 说明“能不能跨过去”。Rust 的线程 API 同时检查这两件事。

但下一个问题很快会出现:如果不是简单的几个任务,而是十万个元素呢?

难道一个元素创建一条线程?

假设要对一个很大的数组执行昂贵计算:

let output: Vec<_> = values
    .iter()
    .map(|value| expensive_compute(*value))
    .collect();

这些计算彼此独立,很适合并行。但如果给每个元素都 thread::spawn,事情会迅速失控。

操作系统线程有自己的栈,也需要内核参与创建和调度。线程不是一个可以无限创建的廉价任务单位。十万个元素不应该对应十万条线程。

我们真正需要的是:

  1. 只创建少量工作线程
  2. 把大量小任务提交给它们
  3. 哪条线程空闲,就继续拿任务执行

这就是线程池存在的理由。

Rust 标准库没有提供通用线程池。对于可以拆分的 CPU 密集型计算,Rayon 是很自然的选择:

use rayon::prelude::*;

fn main() {
    let values: Vec<u64> = (0..100_000).collect();

    let sum: u64 = values
        .par_iter()
        .map(|value| expensive_compute(*value))
        .sum();

    println!("{sum}");
}

fn expensive_compute(value: u64) -> u64 {
    value * value
}

表面上只是把 iter() 换成了 par_iter(),背后却不是“每个元素一条线程”。

Rayon 把计算拆成任务,放到固定大小的工作线程池中。某条工作线程完成自己的任务后,可以从其他线程那里窃取尚未执行的任务。这种工作窃取让负载不均匀的计算也能比较自然地分配多个 CPU 核心上。

par_iter() 不是普通迭代器上的一个并行开关。它来自 Rayon 的并行迭代器 trait,mapfilterreducesum 等操作都需要满足并行执行的语义。

现在可以看到 thread::scope 和 Rayon 的差别:

  • 固定几个彼此不同的任务,可以使用 scoped thread
  • 大量同构、可拆分的数据计算,更适合 Rayon

它们都能形成 Fork-Join,但处在不同抽象层级。前者让我们管理几条线程,后者让我们描述一批并行计算。

如果任务大部分时间都在等待

Rayon 解决了大量 CPU 计算的问题,但并发程序不只有计算。

假设一个服务要同时发出一千个网络请求。每个请求真正使用 CPU 的时间可能很短,大部分时间都在等待远端响应。

当然可以为每个请求创建一条线程,但这些线程多数时候只是阻塞在那里。线程栈和调度成本仍然存在,CPU 却没有在做对应数量的工作。

异步任务解决的是另一类问题。

use std::time::Duration;

#[tokio::main]
async fn main() {
    let mut handles = Vec::new();

    for id in 0..1_000 {
        handles.push(tokio::spawn(async move {
            tokio::time::sleep(Duration::from_millis(100)).await;
            id
        }));
    }

    for handle in handles {
        println!("{}", handle.await.unwrap());
    }
}

tokio::spawn 创建的是异步任务,不是“一次调用就创建一条操作系统线程”。

多个任务由 runtime 调度。当任务运行到 .await,并且等待的操作还没有准备好时,它可以把执行机会交还给 runtime。工作线程不必陪着这个任务一起等待,可以继续推进其他已经就绪的任务。

这种调度是协作式的。runtime 不会在任意一行代码之间强行暂停任务;任务通常要运行到一次尚未就绪的 .await,才会把执行权交出去。因此,异步任务是否足够“轻”,不仅取决于创建成本,也取决于它是否愿意及时让出工作线程。

所以异步特别适合大量等待型工作:

  • 网络连接
  • 数据库请求
  • 定时器
  • 支持异步接口的文件或消息系统

但异步不是更轻量的并行计算。

tokio::spawn(async {
    // 长时间占用 CPU,中间没有 await
    expensive_cpu_work();
});

如果这段计算长时间不让出执行权,它会占住 runtime 的工作线程,其他异步任务也无法得到调度。

阻塞 API 可以通过 spawn_blocking 移出异步工作线程。但持续、大量的 CPU 密集型任务仍然应该控制并发度,并优先考虑 Rayon 或专用线程池。

场景更接近的执行模型
少量、生命周期独立的任务thread::spawn
需要借用当前作用域数据的固定任务thread::scope
大量可拆分的 CPU 计算Rayon
大量 I/O 等待async runtime

任务如何交换数据

任务真正跑起来以后,另一个问题才刚刚开始:

它们之间怎么交换数据?

JoinHandle 只适合在汇合点拿结果

前面的 Fork-Join 已经提供了一种返回结果的方法:

let handle = thread::spawn(calculate_left);
let result = handle.join().unwrap();

如果子任务只返回一次结果,而且主任务本来就要等待它结束,JoinHandle 已经足够。

但有些任务不是算完一次就退出。

例如:

  • 工作线程不断接收新任务
  • 下载线程持续汇报进度
  • 日志线程接收其他线程产生的日志
  • 一个处理阶段不断把结果交给下一个阶段

这种关系不是“结束时返回一个值”,而是执行过程中持续通信。

于是我们需要通道(channel)。

通道:把数据交给另一个任务

先从一条消息开始:

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

fn main() {
    let (tx, rx) = mpsc::channel();

    thread::spawn(move || {
        let result = String::from("finished");
        tx.send(result).unwrap();
    });

    let result = rx.recv().unwrap();
    println!("{result}");
}

tx 是发送端,rx 是接收端。

mpsc 是 multi-producer, single-consumer:可以有多个生产者,但只有一个消费者。Sender 可以克隆,Receiver 不能用同样的方式克隆成多个竞争消费者。

比“线程安全队列”更值得注意的是 send 的签名:

pub fn send(&self, t: T) -> Result<(), SendError<T>>

参数是 T,不是 &T。发送一个值,通常意味着把它的所有权移动到通道中:

let message = String::from("hello");
tx.send(message).unwrap();

println!("{message}"); // message 已经被移动

发送者不再拥有这份数据,接收者取到以后成为新的拥有者。

这是不同于共享状态的另一种并发理念:“Don’t communicate by sharing memory, share memory by communicating.”

多个生产者,以及通道如何结束

现在让多个线程同时发送结果:

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

fn main() {
    let (tx, rx) = mpsc::channel();

    for id in 0..4 {
        let tx = tx.clone();

        thread::spawn(move || {
            tx.send(format!("result from worker {id}")).unwrap();
        });
    }

    drop(tx);

    for message in rx {
        println!("{message}");
    }
}

循环里每个线程拿到一个克隆的 Sender。克隆的是发送能力,不是把整个队列复制一份。

这里的 drop(tx) 很关键。

for message in rx 会持续等待后续消息。只有所有发送端都被释放,接收端才知道“以后不会再有消息了”,随后结束迭代。

循环外最初创建的 tx 如果一直活着,即使四个工作线程都已经退出,接收端仍会认为还有发送者存在,于是继续等待。

关闭通道不是额外发送一个特殊的结束标记。发送端的生命周期本身就是协议的一部分。

无界队列把压力藏进了内存

mpsc::channel() 创建的是无界通道。

这里的“无界”不是说机器真的有无限内存,而是发送方不会因为缓冲区满而等待。生产者可以持续发送,尚未消费的消息不断留在内存里。

当生产速度长期高于消费速度时,队列只会越来越长:

producer: 10000 msg/s
consumer:  1000 msg/s

每秒积压 9000 条消息

程序看起来没有阻塞,延迟和内存占用却在悄悄增长。无界队列常常只是把系统压力从发送方转移到了内存。

如果系统能够接受的积压是有限的,就应该把这个限制表达出来:

use std::sync::mpsc;

fn main() {
    let (tx, rx) = mpsc::sync_channel(100);

    // 队列满以后,send 会等待消费者腾出空间
    tx.send("job").unwrap();

    assert_eq!(rx.recv().unwrap(), "job");
}

sync_channel(100) 最多缓存 100 条尚未接收的消息。缓冲区满时,发送方会阻塞,直到消费者拿走一条消息。

这是一种反馈:

消费者已经处理不过来了,生产者不能继续无限制造积压。

这就是背压。

把容量设为 0 还会得到 rendezvous channel。每次发送都必须等到某个接收操作与它配对,数据不在队列中停留。

这里有一个容易混淆的术语。标准库文档把 mpsc::channel() 称为 asynchronous channel,意思是它的发送端不等待缓冲区,不是说它属于 Rust 的 async/await 体系。rx.recv() 仍然会阻塞当前线程。

如果通信双方是 Tokio 异步任务,通常应该使用 tokio::sync::mpsc,并通过 .send(...).await.recv().await 等待容量和消息,而不是在 runtime 工作线程上调用一个长期阻塞的 recv()

通道可以把程序连成流水线

通道不只是“线程 A 给线程 B 发一个数”。它还可以把处理过程拆成多个阶段。

例如:

读取路径 ──> 解析文件 ──> 汇总结果

每个阶段只关心自己的输入和输出:

use std::sync::mpsc::sync_channel;
use std::thread;

fn parse(path: &str) -> usize {
    // 用字符串长度代替真实解析结果
    path.len()
}

fn main() {
    let files = vec!["access.log", "error.log", "app.log"];
    let (path_tx, path_rx) = sync_channel::<String>(8);
    let (count_tx, count_rx) = sync_channel::<usize>(8);

    thread::scope(|scope| {
        scope.spawn(move || {
            for path in files {
                path_tx.send(path.to_owned()).unwrap();
            }
        });

        scope.spawn(move || {
            for path in path_rx {
                let count = parse(&path);
                count_tx.send(count).unwrap();
            }
        });

        let total: usize = count_rx.into_iter().sum();
        println!("{total}");
    });
}

第一条通道限制待解析文件的数量,第二条限制等待汇总的结果数量。任意阶段变慢,压力都会沿着有界通道往上游传递,而不是一直吃内存。

每个阶段也有自然的关闭方式:

  • 路径生产线程退出,path_tx 被释放
  • 解析线程读完 path_rx 后退出,count_tx 被释放
  • 汇总端读完 count_rx,得到最终结果

流水线很适合所有权能够沿单一方向流动的数据,但并不是所有状态都能这样传下去。

数据无法转移时:共享状态

有些数据确实需要共同访问

设想一个服务中的全局任务表:

  • 请求处理任务要创建记录
  • 后台工作线程要更新进度
  • 查询接口要读取当前状态

这份数据不能简单地从 A 移动到 B,因为 A、B 和 C 都还要继续访问它。

最直接的共享尝试通常会失败:

use std::thread;

fn main() {
    let mut counter = 0;

    for _ in 0..10 {
        thread::spawn(|| {
            counter += 1;
        });
    }

    println!("{counter}");
}

多个线程不能随意持有指向同一个局部变量的可变引用。即使生命周期问题解决了,同时写入同一块内存仍然会形成数据竞争。

要共享一个值,先要解决所有权问题。

Arc 解决共同所有权

Arc<T> 是原子引用计数指针。每次克隆 Arc,都会得到一个指向同一份数据的新拥有者:

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

fn main() {
    let message = Arc::new(String::from("hello"));
    let mut handles = Vec::new();

    for _ in 0..4 {
        let message = Arc::clone(&message);

        handles.push(thread::spawn(move || {
            println!("{message}");
        }));
    }

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

最后一个 Arc 被释放时,内部数据才会被释放。

Arc 只解决“这份数据由谁拥有”,不等于“任何人都能随便修改”。

多个 Arc 指向同一个值时,我们拿到的仍然是共享访问。要修改内部状态,还需要一种能够协调可变访问的机制。

Mutex 解决互斥访问

计数器可以写成:

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

fn main() {
    let counter = Arc::new(Mutex::new(0));
    let mut handles = Vec::new();

    for _ in 0..10 {
        let counter = Arc::clone(&counter);

        handles.push(thread::spawn(move || {
            let mut value = counter.lock().unwrap();
            *value += 1;
        }));
    }

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

    println!("{}", *counter.lock().unwrap());
}
  • Arc<T> 让多条线程共同拥有同一个值
  • Mutex<T> 保证同一时刻只有一个线程能可变访问内部数据
  • MutexGuard 代表当前已经持有锁,离开作用域时自动解锁

Arc 解决所有权,Mutex 解决互斥访问。

如果共享数据本身活得足够久,而且线程都被限制在 thread::scope 中,有时可以直接借用一个 Mutex<T>,不一定需要 Arc。反过来,只读共享数据可能只需要 Arc<T>,不需要 Mutex

先分清问题,类型才不会一层层机械地套上去。

读多写少就一定用 RwLock 吗?

Mutex 同一时刻只允许一个守卫存在,即使大家都只想读取。

RwLock 允许多个读者同时持有读锁,但写者仍然需要独占:

use std::sync::RwLock;

fn main() {
    let config = RwLock::new(String::from("v1"));

    {
        let first = config.read().unwrap();
        let second = config.read().unwrap();
        println!("{first}, {second}");
    }

    *config.write().unwrap() = String::from("v2");
}

它适合读操作明显多于写操作,并且读锁会持有一段有意义时间的场景。

RwLock 不会天然比 Mutex 快。它需要维护更复杂的状态,读写竞争和具体平台实现也会影响结果。一次非常短的临界区,换成 RwLock 未必能得到收益。

“读多写少”是开始考虑它的理由,不是跳过基准测试的理由。

一个计数器需要整把锁吗?

如果共享状态只是一个独立计数器,用 Mutex<usize> 当然可以:

let mut count = counter.lock().unwrap();
*count += 1;

但这类极小粒度操作可以直接使用原子类型:

use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::thread;

fn main() {
    let counter = Arc::new(AtomicUsize::new(0));
    let mut handles = Vec::new();

    for _ in 0..10 {
        let counter = Arc::clone(&counter);

        handles.push(thread::spawn(move || {
            counter.fetch_add(1, Ordering::Relaxed);
        }));
    }

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

    println!("{}", counter.load(Ordering::Relaxed));
}

这里使用 Relaxed,是因为我们只关心计数器自身的原子性,不依赖它来发布或同步其他内存状态。

Ordering 不是从“快”到“慢”的性能档位。它描述的是当前原子操作与其他内存访问之间需要建立怎样的可见性关系。只有在算法确实不依赖额外同步时,Relaxed 才是正确选择。

原子类型适合:

  • 计数器
  • 标志位
  • 简单状态机
  • 统计指标

它不适合把复杂业务状态拆成一堆彼此独立的原子变量。

struct Account {
    balance: i64,
    version: u64,
}

如果 balanceversion 必须作为一个整体变化,把它们分别放进两个 Atomic,读者可能看到新余额和旧版本的组合。

此时一把保护整个 AccountMutex 更容易维持不变量,也更容易读懂。

无锁不是自动获得的性能勋章。复杂状态的正确性,往往比少一次加锁更重要。

锁真正麻烦的地方

Arc<Mutex<T>> 能通过编译,不代表锁就用对了。

最常见的问题是临界区太大:

let mut data = state.lock().unwrap();

let response = slow_network_call();
data.push(response);

网络调用期间,锁一直被当前线程持有。其他线程即使只想快速更新一条数据,也只能等待。

更合适的写法是先完成不需要保护的工作:

let response = slow_network_call();

let mut data = state.lock().unwrap();
data.push(response);

持锁期间尽量只做必须受保护的内存操作

多把锁还会带来死锁:

线程 A:先锁 account,再锁 order
线程 B:先锁 order,再锁 account

如果 A 拿着 accountorder,B 拿着 orderaccount,两边都不会继续。

常见的控制方式包括:

  • 统一加锁顺序
  • 减少同时持有多把锁
  • 把需要共同维护不变量的数据放进同一把锁
  • 让某个任务独占状态,其他任务通过通道向它发送命令

标准库的 Mutex 还有 poisoning。

如果线程持锁期间 panic,内部状态可能只改了一半。之后的 lock() 会通过 PoisonError 提醒调用者:互斥锁仍然能工作,但里面的数据可能已经不满足原来的不变量。

所以 lock().unwrap() 不是所有场景下都理所当然。它表达的是一种策略:一旦状态可能不一致,当前调用者也选择继续 panic。

一个融合的例子

Fork-Join、通道和共享状态不是三选一。

一个图片处理服务可能同时使用它们:

                     ┌── async task:接收请求


              有界 async 通道


          Rayon / 专用线程池:处理图片

          ┌──────────┴──────────┐
          ▼                     ▼
通道返回处理结果         Mutex 保存任务状态


                         Atomic 记录指标

这里至少有三层不同的问题:

  1. 网络连接如何被高效调度
  2. CPU 工作如何限制并行度并利用多核
  3. 结果和状态如何在不同任务之间流动

Tokio 不会因为能 spawn 就替代 Rayon,通道也不会因为减少共享状态就让锁彻底消失。它们是不同层次的构建模块。

到底该从哪里开始选?

可以先看任务的形状:

场景优先考虑
一次性、数量有限的独立任务thread::spawn
固定任务需要借用当前数据thread::scope
大量独立的 CPU 计算Rayon
大量网络或其他 I/O 等待async / Tokio

再看数据如何协作:

场景优先考虑
子任务结束时返回一次结果JoinHandle
生产者持续把数据交给消费者通道
需要限制积压并形成背压有界通道
多个任务必须共同维护复杂状态Arc<Mutex<T>>
读多写少,并且确实存在读并发Arc<RwLock<T>>
独立计数器或简单标志位Atomic
一份状态最好只由一个任务修改独占任务 + 通道
任务彼此独立、最终只需要汇总吗?
    └─ 是:Fork-Join

数据可以在任务之间转移所有权吗?
    └─ 是:通道

多个任务确实必须同时访问同一份状态吗?
    └─ 是:Mutex / RwLock / Atomic

能拆开的数据,先拆开。

能转移的数据,考虑让所有权流动。

确实无法转移、必须共同维护的状态,再引入共享所有权和同步原语。

最后还是所有权

回看这些并发方式,会发现它们一直在回答 Rust 所有权中的几个问题。

Fork-Join 倾向于把数据拆成互不重叠的部分:

let (left, right) = values.split_at_mut(mid);

编译器证明两个任务不会同时修改同一片内存。

通道让所有权从发送者移动到接收者:

tx.send(value);

发送以后,原发送者通常不能再使用 value

共享状态无法把所有权交给单一任务,于是使用:

Arc<Mutex<T>>

Arc 承认有多个拥有者,Mutex 再把运行时访问约束为同一时刻只有一个可变访问者。

Rust 没有发明一套全新的并发问题。任务仍然要拆分、等待和汇合,数据仍然要传递或共享。它做得更彻底的地方,是把这些关系放进所有权、生命周期以及 Send / Sync 约束里,让许多错误在程序运行以前就暴露出来。

面对一个并发问题时,先问两个问题:

任务是什么形状?数据最终归谁所有?

References

  1. Rust 标准库:std::thread
  2. Rust 标准库:std::sync::mpsc
  3. Rust 标准库:Mutex
  4. Rayon
  5. Tokio:CPU-bound tasks and blocking code