Article
共享资源、消息传递与并发控制
问题背景
并发编程解决的是多个任务如何共同推进。但只要任务之间不是完全独立,它们就会遇到另一个问题:如何交换数据,如何协作完成同一件事
常见方式可以分成三类:
- 共享内存:多个任务访问同一份数据,通过 Mutex、RwLock、Atomic 等工具保证一致性
- 共享资源:多个任务使用同一个有限资源,通过 Semaphore 等工具限制并发使用数量
- 消息传递:任务之间发送消息,数据沿着 Channel、Queue 或 Actor Mailbox 流动
这些方式不是互斥关系。很多真实系统会同时使用它们:在局部性能关键路径上使用共享内存,在有限资源入口使用 Semaphore,在模块边界或任务分发处使用消息队列
工具地图
先把常见工具放在一张表里。它们解决的问题不同,不能只按“高级”或“低级”排序
| 工具 | 核心作用 | 典型场景 |
|---|---|---|
| Mutex / 互斥锁 | 同一时间只允许一个线程访问共享数据 | 修改共享变量、写文件 |
| RwLock / 读写锁 | 多读单写 | 配置缓存、读多写少的数据 |
| Atomic / 原子操作 | 对简单变量做无锁修改 | 计数器、状态标记 |
| Semaphore / 信号量 | 限制共享资源的同时使用数量 | 限流、连接池 |
| Condition Variable / 条件变量 | 等待某个条件成立后再继续 | 生产者-消费者 |
| Barrier / 屏障 | 多个线程都到达某点后一起继续 | 并行计算阶段同步 |
| Channel / 通道 | 通过消息传递数据 | Go、Rust、CSP 模型 |
| Queue / 队列 | 任务排队,消费者异步处理 | 后端任务、日志写入 |
| Actor 模型 | 每个对象自己处理消息,不共享状态 | Erlang、Akka、分布式系统 |
共享内存
共享内存是最直接的协作方式。多个线程可以访问同一个对象、缓存、配置或计数器,不需要显式复制数据
它的优势是访问路径短,性能通常较好。缺点是只要共享数据可变,就必须处理同步问题
共享数据可以分成两种情况:
- 只读共享数据:多个任务同时读取通常是安全的
- 可变共享数据:只要存在并发写入,就需要同步控制
例如多个线程同时更新同一个计数器,如果每个线程都执行“读取、加一、写回”,中间任何一步被切换都可能丢失更新
共享状态的三种粒度
共享内存里的同步工具可以按“保护什么”来理解:
- Mutex:保护一整份可变状态
- RwLock:保护读多写少的可变状态
- Atomic:保护一个简单的数值或标记
它们不是递进替代关系,而是适合不同粒度的数据。数据结构越复杂,越需要锁来保护整体不变量;状态越简单,越可能用 Atomic 避免阻塞
Mutex
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..4 {
let counter = Arc::clone(&counter);
let handle = thread::spawn(move || {
for _ in 0..1000 {
let mut value = counter.lock().unwrap();
*value += 1;
}
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
println!("{}", *counter.lock().unwrap());
}这里的 Arc 负责让多个线程共享同一个计数器,Mutex 负责保证同一时刻只有一个线程修改它
Mutex 的关键点是独占。无论线程只是读取,还是准备修改,只要进入 lock() 保护的区域,其他线程都要等待
这让 Mutex 很容易保持一致性,但也会引入新的成本:
- 线程可能阻塞等待锁
- 锁粒度过大会降低并行度
- 多把锁顺序不当可能产生死锁
- 临界区里执行慢操作会放大等待时间
因此,Mutex 不是不能用,而是要控制范围。通常应该让临界区尽量短,只保护真正需要同步的数据修改
RwLock
如果共享数据读多写少,可以考虑 RwLock。它把访问分成读锁和写锁:读锁可以并发持有,写锁必须独占
这适合配置对象、路由表、元数据缓存等场景。大部分时间只是读取,偶尔才需要更新
读锁通过 read() 获取,多个线程可以同时持有读锁。写锁通过 write() 获取,同一时间只能有一个线程持有写锁
只要有写锁存在,就不能再有读锁;只要还有读锁未释放,写锁也必须等待
use std::collections::HashMap;
use std::sync::{Arc, RwLock};
use std::thread;
fn main() {
let config = Arc::new(RwLock::new(HashMap::from([
("timeout".to_string(), "30s".to_string()),
("region".to_string(), "ap-east".to_string()),
])));
let mut readers = Vec::new();
for _ in 0..3 {
let config = Arc::clone(&config);
readers.push(thread::spawn(move || {
let guard = config.read().unwrap();
println!("region = {:?}", guard.get("region"));
}));
}
for reader in readers {
reader.join().unwrap();
}
let writer_config = Arc::clone(&config);
let writer = thread::spawn(move || {
let mut guard = writer_config.write().unwrap();
guard.insert("timeout".to_string(), "60s".to_string());
});
writer.join().unwrap();
}上面的读线程只查看配置,因此拿读锁即可。多个读线程可以同时读取 HashMap,不会互相阻塞
更新配置时必须拿写锁。写锁会等待所有读锁释放,并阻止新的读锁进入,避免读取过程中数据被同时修改
RwLock 的收益来自“读很多,写很少”。如果写入很频繁,读写之间会反复互相等待,它可能并不比 Mutex 更合适
Atomic
如果共享状态只是一个简单数值,可以考虑 Atomic。原子操作由 CPU 指令和语言内存模型提供保证,适合计数器、状态标记、引用计数等小型数据
常见原子操作包括:
load:原子读取当前值store:原子写入新值swap:原子替换,并返回旧值compare_exchange:比较当前值,匹配时才替换fetch_add/fetch_sub:原子加减,并返回旧值fetch_and/fetch_or/fetch_xor:原子位运算fetch_max/fetch_min:原子更新最大值或最小值
use std::sync::atomic::{AtomicUsize, Ordering};
static REQUESTS: AtomicUsize = AtomicUsize::new(0);
fn read_requests() -> usize {
REQUESTS.load(Ordering::Relaxed)
}
fn record_request() {
REQUESTS.fetch_add(1, Ordering::Relaxed);
}多个线程可以同时对同一个原子变量执行 load。原子读取不像 Mutex 那样互斥,也不会阻塞其他读取者
原子读取也可以和原子写入并发发生。例如一个线程执行 fetch_add(1),另一个线程同时执行 load(),这是允许的,也不会形成 data race
读线程可能看到写入前的旧值,也可能看到写入后的新值。但它不会读到“半更新”的坏值,例如乱码、撕裂值或不完整状态
例如 REQUESTS 当前是 10,另一个线程正在执行 fetch_add(1)。并发读取可能得到 10,也可能得到 11,但不会得到非法中间值
把三者放在一起看,区别会更清楚:
| 操作 | 多个读能否同时发生 | 读写能否同时发生 |
|---|---|---|
Mutex::lock() | 不可以 | 不可以 |
RwLock::read() | 可以 | 不可以 |
Atomic::load() | 可以 | 可以 |
Atomic 不是靠“让其他线程等待”来保证安全,而是保证单次读写不可撕裂,并受内存序规则约束
因此,Atomic 的 load-load、load-store、load-fetch_add 都可以并发。它们保证每次读写是合法的,但不保证所有线程读到同一个值
Atomic 不是“更轻的锁”的万能替代品。它适合表达简单状态变化,但复杂不变量通常仍然需要 Mutex、RwLock 或更高层的并发模型
Condition Variable
条件变量用于“等待某个条件成立”。它通常和 Mutex 搭配使用:锁保护共享状态,条件变量负责让线程睡眠和唤醒
生产者-消费者就是典型场景。消费者发现队列为空时,不应该一直循环检查,而是等待条件变量。生产者放入新任务后,再通知消费者继续处理
条件变量容易出错的地方是“唤醒不等于条件一定成立”。被唤醒后仍然要重新检查条件,所以通常会写成 while 循环,而不是 if
下面用 C++ 写一个生产者-消费者模型。mutex 保护共享队列和结束标记,condition_variable 负责让消费者在队列为空时睡眠
#include <condition_variable>
#include <iostream>
#include <mutex>
#include <queue>
#include <thread>
#include <vector>
std::mutex mutex;
std::condition_variable condition;
std::queue<int> tasks;
bool done = false;
void producer(int id) {
for (int i = 0; i < 5; ++i) {
{
std::lock_guard<std::mutex> lock(mutex);
tasks.push(id * 100 + i);
}
condition.notify_one();
}
}
void consumer(int id) {
while (true) {
std::unique_lock<std::mutex> lock(mutex);
while (tasks.empty() && !done) {
condition.wait(lock);
}
if (tasks.empty() && done) {
break;
}
int task = tasks.front();
tasks.pop();
lock.unlock();
std::cout << "consumer " << id << " handles " << task << '\n';
}
}
int main() {
std::vector<std::thread> producers;
std::vector<std::thread> consumers;
for (int i = 0; i < 2; ++i) {
producers.emplace_back(producer, i);
}
for (int i = 0; i < 2; ++i) {
consumers.emplace_back(consumer, i);
}
for (auto& thread : producers) {
thread.join();
}
{
std::lock_guard<std::mutex> lock(mutex);
done = true;
}
condition.notify_all();
for (auto& thread : consumers) {
thread.join();
}
}消费者拿到锁后先检查队列。如果队列为空并且生产还没结束,就调用 wait() 进入睡眠。wait() 会临时释放锁,让生产者有机会拿锁并写入任务
生产者把任务放入队列后调用 notify_one(),唤醒一个等待中的消费者。所有生产者结束后,主线程把 done 改成 true,再用 notify_all() 唤醒剩余消费者退出
这里必须使用 while 重新检查条件,因为线程可能被虚假唤醒,也可能被唤醒后发现任务已经被其他消费者取走。条件变量只负责通知,不负责保证条件仍然成立
Barrier
屏障用于阶段同步。多个线程各自执行一段工作,只有所有线程都到达同一个同步点后,才能一起进入下一阶段
它常见于并行计算。例如把矩阵分块给多个线程处理,每个线程完成本轮计算后,在屏障处等待。所有线程到齐后,再进入下一轮
Barrier 不保护数据,也不传递消息。它解决的是“大家必须一起进入下一阶段”的协调问题
下面是一个 Rust 示例。每个线程先完成第一阶段,然后在 barrier.wait() 处等待。只有 3 个线程都到达屏障后,它们才会继续执行第二阶段
use std::sync::{Arc, Barrier};
use std::thread;
use std::time::Duration;
fn main() {
let barrier = Arc::new(Barrier::new(3));
let mut handles = Vec::new();
for worker_id in 0..3 {
let barrier = Arc::clone(&barrier);
handles.push(thread::spawn(move || {
println!("worker {worker_id} starts phase 1");
thread::sleep(Duration::from_millis(100 * worker_id));
println!("worker {worker_id} waits at barrier");
barrier.wait();
println!("worker {worker_id} starts phase 2");
}));
}
for handle in handles {
handle.join().unwrap();
}
}这个例子里,线程 0 会最先到达屏障,但它不能立刻进入第二阶段。它必须等线程 1 和线程 2 也执行到 wait(),然后三个线程一起继续
Barrier::new(3) 里的 3 表示参与同步的线程数量。如果少于 3 个线程到达屏障,已经到达的线程会一直等待
所以 Barrier 适合阶段式任务,不适合保护共享数据。如果第二阶段要读取第一阶段共同生成的数据,数据本身仍然需要 Mutex、RwLock 或其他同步方式保护
共享资源
共享资源不是一块共享内存,而是多个任务都要使用的有限能力。数据库连接、GPU 推理服务、文件句柄、网络带宽和远端 API 都属于这类问题
这类问题的核心不是“数据会不会被改坏”,而是“同时使用的人太多,资源会不会扛不住”
Semaphore
Semaphore 不解决“共享数据会不会被改坏”的问题,而是解决“共享资源会不会被同时用爆”的问题
这里的容量不是内存大小,也不是 Vec 或 HashMap 能存多少元素,而是并发许可数量
例如 Semaphore::new(3) 表示同一时间最多允许 3 个任务进入被保护的资源区域。第 4 个任务必须等待,直到前面的任务释放许可
共享资源可以是数据库连接、GPU 推理服务、文件句柄、网络带宽、远端 API 或线程池。重点不是数据一致性,而是资源承载能力
| 共享资源 | 容量含义 |
|---|---|
| 数据库连接池 | 最多同时使用多少个连接 |
| GPU 推理服务 | 最多同时运行多少个推理任务 |
| 文件下载 | 最多同时进行多少个下载 |
| 远端 API | 最多同时发出多少个请求 |
| 打印机 | 最多同时处理多少个打印任务 |
use std::sync::Arc;
use tokio::sync::Semaphore;
async fn fetch_with_limit(semaphore: Arc<Semaphore>) {
let _permit = semaphore.acquire().await.unwrap();
// 同一时间最多只有固定数量的任务能进入这里
fetch_remote_data().await;
}实际项目里常见的是 tokio::sync::Semaphore。它适合限制异步任务对有限资源的并发使用数量
共享内存和共享资源可以同时出现。例如一个后端推理服务既要修改任务状态表,又要使用 GPU 跑模型:
use std::collections::HashMap;
use std::sync::Mutex;
use tokio::sync::Semaphore;
struct JobStatus;
struct AppState {
jobs: Mutex<HashMap<String, JobStatus>>,
gpu_slots: Semaphore,
}这里的 jobs 是共享数据,怕被多个线程同时修改出错,所以用 Mutex 保护。gpu_slots 表示 GPU 的并发容量,怕同时推理任务太多,所以用 Semaphore 限制
一句话区分:Mutex、RwLock、Atomic 保护共享数据的一致性;Semaphore 限制共享资源的同时使用数量
消息传递
消息传递的核心思想是:不要让多个任务同时修改同一份数据,而是让数据沿着明确的方向流动
也可以换一种说法:不是通过共享内存来通信,而是通过通信来转移数据或表达意图
在这种模型下,发送方把消息交给 Channel 或 Queue,接收方从中取出消息并处理。数据修改权集中在接收方,发送方不再直接操作接收方的内部状态
典型场景包括:
- 后台任务分发
- 日志异步写入
- Actor 模型
- 事件发布订阅
- 消息队列消费
Channel
Channel 是进程内或运行时内的消息通道。它通常由发送端和接收端组成,发送端负责投递消息,接收端负责接收消息
Rust 标准库的 mpsc 就是一个多生产者、单消费者 Channel。mpsc 表示 multiple producer, single consumer
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
let worker = thread::spawn(move || {
while let Ok(message) = rx.recv() {
println!("worker received: {message}");
}
});
tx.send("compile").unwrap();
tx.send("test").unwrap();
tx.send("deploy").unwrap();
drop(tx);
worker.join().unwrap();
}这个例子里,主线程不直接修改 worker 的内部状态,而是把命令发送给 worker。worker 按顺序接收并处理消息
Channel 的好处是边界清晰。任务之间的协作被建模成消息流,谁发送、谁接收、数据往哪里走,都比较容易看出来
消息队列
消息队列可以理解为带缓冲能力的消息通道。生产者把消息放入队列,消费者从队列中取出消息
队列的意义不只是“传数据”,还包括削峰和解耦:
- 生产者和消费者不必同速运行
- 消费者暂时变慢时,队列可以缓冲一部分任务
- 可以增加消费者数量,提高处理吞吐
- 生产逻辑和消费逻辑可以分开部署或分开演进
在进程内,队列可能只是一个内存结构。在分布式系统中,队列通常由 Kafka、RabbitMQ、Redis Stream、SQS 等中间件提供
Actor 模型
Actor 模型可以看作消息传递的一种组织方式。每个 Actor 拥有自己的内部状态,只能通过接收消息来修改状态
外部任务不能直接改 Actor 的字段,只能给它发消息。Actor 在自己的执行上下文中串行处理消息,因此可以避免很多共享可变状态问题
Erlang、Akka 和很多分布式系统都使用类似思想。它适合状态边界清晰、模块之间通过消息协作的场景
Actor 模型的代价是调用路径不再像普通函数那样直接。消息可能排队、失败、超时,也需要考虑 mailbox 积压和错误恢复
生产者-消费者模式
生产者-消费者模式是消息队列最常见的使用方式。生产者负责生成任务,消费者负责取出并处理任务
下面是一个多生产者、多消费者的简化实现。由于 Rust 标准库的 Receiver 不能直接 clone,这里用 Arc<Mutex<Receiver>> 让多个消费者共享接收端
use std::sync::{mpsc, Arc, Mutex};
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel::<i32>();
let rx = Arc::new(Mutex::new(rx));
let mut producers = Vec::new();
for producer_id in 0..2 {
let tx = tx.clone();
producers.push(thread::spawn(move || {
for task_id in 0..5 {
let task = producer_id * 100 + task_id;
tx.send(task).unwrap();
}
}));
}
drop(tx);
let mut consumers = Vec::new();
for consumer_id in 0..2 {
let rx = Arc::clone(&rx);
consumers.push(thread::spawn(move || {
loop {
let message = {
let receiver = rx.lock().unwrap();
receiver.recv()
};
match message {
Ok(task) => {
println!("consumer {consumer_id} handles {task}");
}
Err(_) => break,
}
}
}));
}
for producer in producers {
producer.join().unwrap();
}
for consumer in consumers {
consumer.join().unwrap();
}
}这个例子仍然用了锁,但锁保护的是接收端,而不是业务数据本身。业务任务通过 Channel 流动,消费者只处理自己取到的消息
如果需要真正的多生产者、多消费者 Channel,通常会选择 crossbeam-channel、flume 或异步运行时提供的 Channel,而不是手动把标准库 Receiver 包在锁里
三种方式对比
| 对比维度 | 共享内存 | 共享资源 | 消息传递 |
|---|---|---|---|
| 关注对象 | 同一份数据 | 有限资源能力 | 流动的消息 |
| 主要问题 | 数据会不会被并发读写破坏 | 同时使用者太多会不会超载 | 任务之间如何解耦协作 |
| 常见工具 | Mutex、RwLock、Atomic、Condvar、Barrier | Semaphore | Channel、Queue、Actor |
| 主要风险 | 数据竞争、死锁、锁竞争 | 许可泄漏、容量设置不合理、请求排队 | 队列积压、消息丢失、背压不足 |
| 适合场景 | 读多写少、局部共享、性能关键路径 | 连接池、限流、GPU、远端 API | 任务分发、模块解耦、异步流水线 |
选择时可以先问两个问题:你要协调的是一份数据,还是一个有限资源?这份工作是否天然应该交给某个任务独立处理?
如果多个任务必须访问同一份数据,优先考虑共享内存工具。它们解决的是共享数据一致性问题
如果多个任务争用的是数据库连接、GPU、API 配额或网络带宽,优先考虑 Semaphore。它解决的是共享资源容量问题
如果一份工作天然属于某个处理者,优先考虑消息传递。让数据沿着 Channel、Queue 或 Actor Mailbox 流动,可以减少共享可变状态
总结
共享内存强调共同访问同一份数据,因此需要 Mutex、RwLock、Atomic、Condition Variable 或 Barrier 等工具协调访问
共享资源强调共同使用有限资源,因此需要 Semaphore 这类工具限制并发容量
消息传递强调数据沿着明确方向流动,因此更适合 Channel、Queue、Actor、异步处理和生产者-消费者模型
参考文献
- Rust 圣经:使用消息传递在线程间传送数据
- Rust 圣经:共享状态并发
- Go by Example: Channels
- 如果不想在本机配置环境,C++ 代码片段可以放到 C++ Playground 运行,Rust 代码片段可以放到 Rust Playground 运行