共享资源、消息传递与并发控制

问题背景

并发编程解决的是多个任务如何共同推进。但只要任务之间不是完全独立,它们就会遇到另一个问题:如何交换数据,如何协作完成同一件事

常见方式可以分成三类:

  • 共享内存:多个任务访问同一份数据,通过 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 是最直接的锁。同一时间只允许一个线程进入临界区,适合保护会被多个线程修改的共享数据

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。它把访问分成读锁和写锁:读锁可以并发持有,写锁必须独占

这适合配置对象、路由表、元数据缓存等场景。大部分时间只是读取,偶尔才需要更新

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 指令和语言内存模型提供保证,适合计数器、状态标记、引用计数等小型数据

atom

常见原子操作包括:

  • 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-loadload-storeload-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

Semaphore 不解决“共享数据会不会被改坏”的问题,而是解决“共享资源会不会被同时用爆”的问题

这里的容量不是内存大小,也不是 VecHashMap 能存多少元素,而是并发许可数量

例如 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-channelflume 或异步运行时提供的 Channel,而不是手动把标准库 Receiver 包在锁里

三种方式对比

对比维度共享内存共享资源消息传递
关注对象同一份数据有限资源能力流动的消息
主要问题数据会不会被并发读写破坏同时使用者太多会不会超载任务之间如何解耦协作
常见工具Mutex、RwLock、Atomic、Condvar、BarrierSemaphoreChannel、Queue、Actor
主要风险数据竞争、死锁、锁竞争许可泄漏、容量设置不合理、请求排队队列积压、消息丢失、背压不足
适合场景读多写少、局部共享、性能关键路径连接池、限流、GPU、远端 API任务分发、模块解耦、异步流水线

选择时可以先问两个问题:你要协调的是一份数据,还是一个有限资源?这份工作是否天然应该交给某个任务独立处理?

如果多个任务必须访问同一份数据,优先考虑共享内存工具。它们解决的是共享数据一致性问题

如果多个任务争用的是数据库连接、GPU、API 配额或网络带宽,优先考虑 Semaphore。它解决的是共享资源容量问题

如果一份工作天然属于某个处理者,优先考虑消息传递。让数据沿着 Channel、Queue 或 Actor Mailbox 流动,可以减少共享可变状态

总结

共享内存强调共同访问同一份数据,因此需要 Mutex、RwLock、Atomic、Condition Variable 或 Barrier 等工具协调访问

共享资源强调共同使用有限资源,因此需要 Semaphore 这类工具限制并发容量

消息传递强调数据沿着明确方向流动,因此更适合 Channel、Queue、Actor、异步处理和生产者-消费者模型

参考文献