多线程并发编程
一、预备概念
1.1 并发(Concurrent)与并行(Parallel)
支持多个任务同时存在的系统,称为并发系统;
支持多个任务同时执行的系统,称为并行系统。
1.2 线程(Thread)、进程(Process)与编程语言
区分逻硬线程/逻辑CPU与软件线程/OS线程。
前者是将CPU核心暴露为多个逻辑处理器得到的结果,如8核16线程;后者是操作系统上的概念,即一条执行的指令流。操作系统通过线程调度将线程分配给CPU核心。
运行程序时,操作系统会创建进程,一个进程内可以有多条线程,他们共享进程内的某些资源。比如,每个线程通常有自己的寄存器、栈与程序计数器,但是共享进程的代码段,堆,地址空间,文件描述符。
编程语言往往有自己的编程模型,可以通过操作系统的接口创建线程,应对并发需求。RUST的编程模型是1:1的,即一个线程对应一个CPU内核线程。
二、多线程
2.1 多线程的风险
- 竞态条件(race conditions):多个线程以非一致性的顺序同时访问数据资源
- 死锁(deadlocks):两个线程都想使用某个资源,但是又需要等待对方释放资源后才能使用,结果无法继续执行
- 其他的由于多线程导致的BUG,难以 复现和解决
2.2 创建线程
thread::spawn创建线程,并返回JoinHandle<T>句柄,可以用于主线管理操控子线程。
use std::thread;
use std::time::Duration;
fn main(){
thread::spawn(||{
for i in 1..10{
println!("hi number {} from the spawned thread!",i);
thread::sleep(Duration::from_millis(1));
}
});
for i in 1..5 {
println!("hi number {} from the main thread!", i);
thread::sleep(Duration::from_millis(1));
}
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
输出:
hi number 1 from the main thread!
hi number 1 from the spawned thread!
hi number 2 from the spawned thread!
hi number 2 from the main thread!
hi number 3 from the spawned thread!
hi number 3 from the main thread!
hi number 4 from the spawned thread!
hi number 4 from the main thread!
hi number 5 from the spawned thread!2
3
4
5
6
7
8
9
- 线程内部代码用闭包执行
- main线程结束,则程序也结束,故其需要比子线程存活更久
thread::sleep休眠当前线程,此时会调度运行其他线程
显然输出结果的顺序不是确定的,取决于操作系统的行为,因此对线程的执行顺序不要依赖。
2.3 阻塞线程
如.join(),返回Result<T,E>:
use std::thread;
use std::time::Duration;
fn main() {
let handle = thread::spawn(|| {
for i in 1..5 {
println!("hi number {} from the spawned thread!", i);
thread::sleep(Duration::from_millis(1));
}
});
handle.join().unwrap();
for i in 1..5 {
println!("hi number {} from the main thread!", i);
thread::sleep(Duration::from_millis(1));
}
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
2.4 在线程闭包中使用move
直接在线程闭包中使用其他线程的数据会有一些问题,比如数据的生命周期和线程使用的冲突。
使用move关键字可以达成所有权的转移。
use std::thread;
fn main() {
let v = vec![1, 2, 3];
let handle = thread::spawn(move || {
println!("Here's a vector: {:?}", v);
});
handle.join().unwrap();
}2
3
4
5
6
7
8
9
10
11
12
2.5 线程的终止
- 当
main线程终止时,所有线程自然终止 - 当线程的代码执行完毕时,线程终止
- 当代码持续执行时,父线程(不是
main)的终止并不会强制结束子线程,此时子线程会运行至跑满一个CPU核心,最终直到main线程的结束
use std::thread;
use std::time::Duration;
fn main() {
// 创建一个线程A
let new_thread = thread::spawn(move || {
// 再创建一个线程B
thread::spawn(move || {
loop {
println!("I am a new thread.");
}
})
});
// 等待新创建的线程执行完成
new_thread.join().unwrap();
println!("Child thread is finish!");
// 睡眠一段时间,看子线程创建的子线程是否还在运行
thread::sleep(Duration::from_millis(100));
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
2.6 多线程的性能
- 创建线程需要时间,并且随着线程增多而增大
- 对于CPU密集型任务,线程数等于核心数即可。对于长期阻塞的任务,虽然也可以增加线程数,但是往往可以通过其他方法解决(如
async/await的并发模型)
2.6.1 多线程的开销
多线程的开销是指,性能并不会随着线程数增大线性增长,在线程数超过一定范围后还可能下降。
CAS即Compare-And-Swap,比较并交换,是CPU提供的一种原子操作,常用于实现无锁并发结构。
以下是一个CAS的Hashmap在多线程的应用:
for i in 0..num_threads {
let ht = Arc::clone(&ht);
let handle = thread::spawn(move || {
for j in 0..adds_per_thread {
let key = thread_rng().gen::<u32>();
let value = thread_rng().gen::<u32>();
ht.set_item(key, value);
}
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
- 无锁并不是没有同步,只不过是将等待锁的过程换成了原子操作的竞争与重试
- 读操作可以线性增长,但是强竞争的写操作不行
- 线程过多,CPU的缓存命中率会下降,多个线程容易竞争一条cpu cache line(和cpu的缓存一致性有关)
- 内存带宽瓶颈
2.7 线程屏障(Barrier)
Barrier
use std::sync::{Arc, Barrier};
use std::thread;
fn main() {
let mut handles = Vec::with_capacity(6);
let barrier = Arc::new(Barrier::new(6));
for _ in 0..6 {
let b = barrier.clone();
handles.push(thread::spawn(move|| {
println!("before wait");
b.wait();
println!("after wait");
}));
}
for handle in handles {
handle.join().unwrap();
}
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
最终得到:
before wait
before wait
before wait
before wait
before wait
before wait
after wait
after wait
after wait
after wait
after wait
after wait2
3
4
5
6
7
8
9
10
11
12
2.8 线程局部变量
所谓线程局部变量,在定义后,所有线程都可以获取相同的初始值,且对应的数据在线程间独立。
2.8.1 标准库thread_local
use std::cell::RefCell;
use std::thread;
thread_local!(static FOO: RefCell<u32> = RefCell::new(1));
FOO.with(|f| {
assert_eq!(*f.borrow(), 1);
*f.borrow_mut() = 2;
});
// 每个线程开始时都会拿到线程局部变量的FOO的初始值
let t = thread::spawn(move|| {
FOO.with(|f| {
assert_eq!(*f.borrow(), 1);
*f.borrow_mut() = 3;
});
});
// 等待线程完成
t.join().unwrap();
// 尽管子线程中修改为了3,我们在这里依然拥有main线程中的局部值:2
FOO.with(|f| {
assert_eq!(*f.borrow(), 2);
});2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
使用thread_local!宏初始化变量,并通过with(|| {})的方式在线程内内访问变量,这里f的类型是初始化类型的不可变引用。FOO的类型为LocalKey<T>
在结构体中使用局部线程变量有两种方式,类型命名空间和结构体字段
- 类型命名空间,通过
impl的方式将线程局部变量绑定到结构体的命名空间中,不依赖实例。
use std::cell::RefCell;
struct Foo;
impl Foo{
thread_local!{
static FOO :RefCell<usize>=RefCell::new(0);
}
}
fn main(){
Foo::FOO.with(|x| {
println!("{:?}",x);
})
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
- 结构体保存 LocalKey 引用
thread_local! {
static FOO: RefCell<usize> = RefCell::new(0);
}
struct Bar {
foo: &'static LocalKey<RefCell<usize>>,
}
// 初始化
let bar = Bar { foo: &FOO };
// 调用
bar.foo.with(|x| ...);2
3
4
5
6
7
8
9
10
11
12
13
两种方式最终调用的方式都是LocalKey<T>::with()。
2.8.2 第三方库thread-local
| 对比 | std::thread_local! | thread_local::ThreadLocal<T> |
|---|---|---|
| 线程结束后 | 通常销毁该线程的TLS 数据 | 容器可以保留数据 |
| 汇总数据 | 需要自行传递结果 | 可以通过迭代器汇总 |
use thread_local::ThreadLocal;
use std::sync::Arc;
use std::cell::Cell;
use std::thread;
let tls = Arc::new(ThreadLocal::new());
let mut v = vec![];
// 创建多个线程
for _ in 0..5 {
let tls2 = tls.clone();
let handle = thread::spawn(move || {
// 将计数器加1
// 请注意,由于线程 ID 在线程退出时会被回收,因此一个线程有可能回收另一个线程的对象
// 这只能在线程退出后发生,因此不会导致任何竞争条件
let cell = tls2.get_or(|| Cell::new(0));
cell.set(cell.get() + 1);
});
v.push(handle);
}
for handle in v {
handle.join().unwrap();
}
// 一旦所有子线程结束,收集它们的线程局部变量中的计数器值,然后进行求和
let tls = Arc::try_unwrap(tls).unwrap();
let total = tls.into_iter().fold(0, |x, y| {
// 打印每个线程局部变量中的计数器值,发现不一定有5个线程,
// 因为一些线程已退出,并且其他线程会回收退出线程的对象
println!("x: {}, y: {}", x, y.get());
x + y.get()
});
// 和为5
assert_eq!(total, 5);2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
注意这里的Cell可以在线程间复用,由库的线程ID管理等机制决定。
2.9 互斥锁与条件变量
use std::thread;
use std::sync::{Arc, Mutex, Condvar};
fn main() {
let pair = Arc::new((Mutex::new(false), Condvar::new()));
let pair2 = pair.clone();
thread::spawn(move|| {
let (lock, cvar) = &*pair2;
let mut started = lock.lock().unwrap();
println!("changing started");
*started = true;
cvar.notify_one();
});
let (lock, cvar) = &*pair;
let mut started = lock.lock().unwrap();
while !*started {
started = cvar.wait(started).unwrap();
}
println!("started changed");
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
Mutex<bool> ├── 锁的状态:锁定 / 未锁定 └── 被保护的数据:started = false
Condvar └── 负责让线程等待以及通知线程
wait(原来的 MutexGuard):
1. 接收原来的 MutexGuard
2. 释放对应的 Mutex
3. 让当前线程进入条件变量等待状态
4. 等待被唤醒
5. 重新获取对应的 Mutex
6. 返回新的 MutexGuard
| 时刻 | 主线程 | 子线程 |
|---|---|---|
| T1 | lock() 成功,持有锁 | 尚未运行 |
| T2 | 检查 started == false | 执行 lock(),但被阻塞 |
| T3 | 调用 wait(),释放锁并等待通知 | 等待获得锁 |
| T4 | 处于等待状态 | lock() 成功,获得锁 |
| T5 | 继续等待 | 修改 started = true |
| T6 | 收到 notify_one(),尝试重新获取锁 | 仍然持有锁 |
| T7 | 重新获得锁,wait() 返回 | 释放锁 |
| T8 | 检查 started == true,继续执行 | 结束 |
主线程获得锁,进入循环,cvar进入等待状态,释放锁,子线程获得锁(Mutex::lock()是阻塞式的,会等待直到锁可获取),修改锁保护的变量,唤醒主线程(notify_one()会随机唤醒一个等待中的线程),获取修改后的变量,循环结束。
2.10 执行一次INIT.call_once
use std::thread;
use std::sync::Once;
static mut VAL: usize = 0;
static INIT: Once = Once::new();
fn main() {
let handle1 = thread::spawn(move || {
INIT.call_once(|| {
unsafe {
VAL = 1;
}
});
});
let handle2 = thread::spawn(move || {
INIT.call_once(|| {
unsafe {
VAL = 2;
}
});
});
handle1.join().unwrap();
handle2.join().unwrap();
println!("{}", unsafe { VAL });
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
这里使用unsafe是因为全局可变变量的设置,可以通过设置static VAL: OnceLock<usize> = OnceLock::new();然后VAL.get_or_init(|| 1);来初始化,从而避免unsafe。
三、线程同步
消息传递可以用于进行线程间数据的共享和传递。该章节主要介绍Rust的标准库channel。
3.1 mpsc:mulpitle producer, single consumer
标准库提供接口std::sync::mpsc,暂时只支持多个发送者,一个接收者(其余实现可以寻求第三方库)。
use std::sync::mpsc;
use std::thread;
fn main() {
// 创建一个消息通道, 返回一个元组:(发送者,接收者)
let (tx, rx) = mpsc::channel();
// 创建线程,并发送消息
thread::spawn(move || {
// 发送一个数字1, send方法返回Result<T,E>,通过unwrap进行快速错误处理
tx.send(1).unwrap();
// 下面代码将报错,因为编译器自动推导出通道传递的值是i32类型,那么Option<i32>类型将产生不匹配错误
// tx.send(Some(1)).unwrap()
});
// 在主线程中接收子线程发送的消息并输出
println!("receive {}", rx.recv().unwrap());
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
- 出于对生命周期的考量,
tx的所有权需要转移进子线程。 send()会返回Result<>T,E,因为发送消息有可能会报错,相应地,在发送者关闭的情况下接收者也会返回错误- 发送和接收的类型是由泛型自动推导得到的,如
mpsc::Sender<i32>和mpsc::Receiver<i32>,此后该通道也只能传递该类型的数据 rx.recv()是阻塞式的
3.1.1 非阻塞接收方法rx.recv()
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
tx.send(1).unwrap();
});
println!("receive {:?}", rx.try_recv());
}2
3
4
5
6
7
8
9
10
11
12
显然,线程创建时间大于println!()的时间,因此rx.try_recv()会返回empty报错。
3.1.2 传输中数据的所有权
没有实现copy特征的变量会在发送后转移所有权
3.1.3 通过for循环接收
use std::sync::mpsc;
use std::thread;
use std::time::Duration;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let vals = vec![
String::from("hi"),
String::from("from"),
String::from("the"),
String::from("thread"),
];
for val in vals {
tx.send(val).unwrap();
thread::sleep(Duration::from_secs(1));
}
});
for received in rx {
println!("Got: {}", received);
}
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
需要指出的是,此处单发单接的情况下,消息遵循FIFO顺序被接收。
3.1.3.1 多个发送者
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
let tx1 = tx.clone();
thread::spawn(move || {
tx.send(String::from("hi from raw tx")).unwrap();
});
thread::spawn(move || {
tx1.send(String::from("hi from cloned tx")).unwrap();
});
for received in rx {
println!("Got: {}", received);
}
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
需要将发送者克隆后传入对应的发送线程;
这里需要知道线程创建顺序、发送顺序和接受顺序无法确认。
3.1.4 同步与异步通道
mpsc::channel创建的通道是异步的,即发送者不会因为未接收而阻塞。
可以用mpsc::sync_channel创建同步通道,此时发送消息是阻塞的,直到消息被接收后才解除。
3.1.4.1 消息缓存
mpsc::synv_channel(N)中的N,代表可以发送N条无阻塞的消息,当消息缓存队列满了之后(可以被接收者消费),后续的消息就按同步消息发送。
3.1.5 通道的关闭
当所有发送者或所有接收者都drop后,通道自动关闭,且在编译期实现。
use std::sync::mpsc;
fn main() {
use std::thread;
let (send, recv) = mpsc::channel();
let num_threads = 3;
for i in 0..num_threads {
let thread_send = send.clone();
thread::spawn(move || {
thread_send.send(i).unwrap();
println!("thread {:?} finished", i);
});
}
// 在这里drop send...
for x in recv {
println!("Got: {}", x);
}
println!("finished iterating");
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
上述代码中原始的send一直没被drop,因此会造成主线程的阻塞。
for循环等价于
loop {
match recv.recv() {
Ok(x) => println!("Got: {}", x),
Err(_) => break,
}
}2
3
4
5
6
当还有sender存在时,接受者会阻塞等待。
3.1.6 传递多种类型的数据
可以通过枚举类型实现,也可以为每种类型单独建立通道。
需要注意的是,通道会按枚举中内存占用最大的成员进行对齐,因此会造成内存上的浪费。
3.1.7 第三方库
crossbeam 和 flume
3.2 锁、Condvar和信号量(共享内存)
相比于消息传递,共享内存实现线程同步的特点有:更加简洁(不等于简单)的实现,以及对更高性能的追求。
3.2.1 互斥锁 Mutex
Mutual exclusion
3.2.1.1 单线程中使用Mutex
use std::sync::Mutex;
fn main() {
// 使用`Mutex`结构体的关联函数创建新的互斥锁实例
let m = Mutex::new(5);
{
// 获取锁,然后deref为`m`的引用
// lock返回的是Result
let mut num = m.lock().unwrap();
*num = 6;
// 锁自动被drop
}
println!("m = {:?}", m);
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
lock()是阻塞式的,且有可能报错,比如持有锁的线程panic时,会返回一个错误lock()返回智能指针MutexGuard<T>:
- 实现了
Deref特征,自动解引用后返回一个引用类型,指向Mutex内部的数据 - 实现了
Drop特征,超出作用域后自动释放锁
3.2.1.2 多线程中使用Mutex
多线程安全的Arc<T>
use std::sync::{Arc, Mutex};
use std::thread;
fn main() {
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());
//这里.lock()自动找到counter内部的Mutex,*再对返回的MutexGuard进行解引用,获得内部数据
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
3.2.1.3 内部可变性
内部可变性的核心思想是,即使只有不可变引用,也可以通过特定类型来安全地修改内部数据。
Rc<T>/RefCell<T>用于单线程内部可变性,Arc<T>/Mutex<T>用于多线程内部可变性。
Arc用于让多线程共享引用,Mutex避免写入竞争,从而安全地修改数据,从而实现多线程的内部可变性。
3.2.1.4 使用Mutex需要注意的问题
使用数据前后必须注意获取锁与释放锁
Mutex也有自己的风险,如造成死锁
3.2.2 死锁
3.2.2.1 单线程死锁
use std::sync::Mutex;
fn main() {
let data = Mutex::new(0);
let d1 = data.lock();
let d2 = data.lock();
} // d1锁在此处释放2
3
4
5
6
7
在d2处死锁
3.2.2.2 多线程死锁
use std::{sync::{Mutex, MutexGuard}, thread};
use std::thread::sleep;
use std::time::Duration;
use lazy_static::lazy_static;
lazy_static! {
static ref MUTEX1: Mutex<i64> = Mutex::new(0);
static ref MUTEX2: Mutex<i64> = Mutex::new(0);
}
fn main() {
// 存放子线程的句柄
let mut children = vec![];
for i_thread in 0..2 {
children.push(thread::spawn(move || {
for _ in 0..1 {
// 线程1
if i_thread % 2 == 0 {
// 锁住MUTEX1
let guard: MutexGuard<i64> = MUTEX1.lock().unwrap();
println!("线程 {} 锁住了MUTEX1,接着准备去锁MUTEX2 !", i_thread);
// 当前线程睡眠一小会儿,等待线程2锁住MUTEX2
sleep(Duration::from_millis(10));
// 去锁MUTEX2
let guard = MUTEX2.lock().unwrap();
// 线程2
} else {
// 锁住MUTEX2
let _guard = MUTEX2.lock().unwrap();
println!("线程 {} 锁住了MUTEX2, 准备去锁MUTEX1", i_thread);
let _guard = MUTEX1.lock().unwrap();
}
}
}));
}
// 等子线程完成
for child in children {
let _ = child.join();
}
println!("死锁没有发生");
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
3.2.2.3 try_lock
显然,为了避免阻塞,try_lock在一次尝试失败后返回错误
use std::{sync::{Mutex, MutexGuard}, thread};
use std::thread::sleep;
use std::time::Duration;
use lazy_static::lazy_static;
lazy_static! {
static ref MUTEX1: Mutex<i64> = Mutex::new(0);
static ref MUTEX2: Mutex<i64> = Mutex::new(0);
}
fn main() {
// 存放子线程的句柄
let mut children = vec![];
for i_thread in 0..2 {
children.push(thread::spawn(move || {
for _ in 0..1 {
// 线程1
if i_thread % 2 == 0 {
// 锁住MUTEX1
let guard: MutexGuard<i64> = MUTEX1.lock().unwrap();
println!("线程 {} 锁住了MUTEX1,接着准备去锁MUTEX2 !", i_thread);
// 当前线程睡眠一小会儿,等待线程2锁住MUTEX2
sleep(Duration::from_millis(10));
// 去锁MUTEX2
let guard = MUTEX2.try_lock();
println!("线程 {} 获取 MUTEX2 锁的结果: {:?}", i_thread, guard);
// 线程2
} else {
// 锁住MUTEX2
let _guard = MUTEX2.lock().unwrap();
println!("线程 {} 锁住了MUTEX2, 准备去锁MUTEX1", i_thread);
sleep(Duration::from_millis(10));
let guard = MUTEX1.try_lock();
println!("线程 {} 获取 MUTEX1 锁的结果: {:?}", i_thread, guard);
}
}
}));
}
// 等子线程完成
for child in children {
let _ = child.join();
}
println!("死锁没有发生");
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
现在这段代码不会死锁了,返回:
线程 0 锁住了MUTEX1,接着准备去锁MUTEX2 !
线程 1 锁住了MUTEX2, 准备去锁MUTEX1
线程 1 获取 MUTEX1 锁的结果: Err("WouldBlock")
线程 0 获取 MUTEX2 锁的结果: Err("WouldBlock")
死锁没有发生2
3
4
5
3.2.3 读写锁RwLock
允许并发读
use std::sync::RwLock;
fn main() {
let lock = RwLock::new(5);
// 同一时间允许多个读
{
let r1 = lock.read().unwrap();
let r2 = lock.read().unwrap();
assert_eq!(*r1, 5);
assert_eq!(*r2, 5);
} // 读锁在此处被drop
// 同一时间只允许一个写
{
let mut w = lock.write().unwrap();
*w += 1;
assert_eq!(*w, 6);
// 以下代码会阻塞发生死锁,因为读和写不允许同时存在
// 写锁w直到该语句块结束才被释放,因此下面的读锁依然处于`w`的作用域中
// let r1 = lock.read();
// println!("{:?}",r1);
}// 写锁在此处被drop
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
- 同时允许多个读,但最多只能有一个写
- 读和写不能同时存在
- 读可以使用
read、try_read,写write、try_write, 在实际项目中,try_xxx会安全的多
问题在于,RwLock的性能往往不如Mutex,实现更复杂,问题也更多,往往开销会更大。 总之,如果你要使用RwLock要确保满足以下两个条件:并发读,且需要对读到的资源进行"长时间"的操作,这里的长时间是持有读锁的时间。
3.2.4 三方库
parking_lot
3.2.5 条件变量Condvar
主要是用于解决读取顺序的问题,常用方法有:wait(),notify_one()
use std::sync::{Arc,Mutex,Condvar};
use std::thread::{spawn,sleep};
use std::time::Duration;
fn main() {
let flag = Arc::new(Mutex::new(false));
let cond = Arc::new(Condvar::new());
let cflag = flag.clone();
let ccond = cond.clone();
//子线程
let hdl = spawn(move || {
let mut lock = cflag.lock().unwrap();
let mut counter = 0;
while counter < 3 {
while !*lock {
// wait方法会接收一个MutexGuard<'a, T>,且它会自动地暂时释放这个锁,使其他线程可以拿到锁并进行数据更新。
// 同时当前线程在此处会被阻塞,直到被其他地方notify后,它会将原本的MutexGuard<'a, T>还给我们,即重新获取到了锁,同时唤醒了此线程。
lock = ccond.wait(lock).unwrap();
}
*lock = false;
counter += 1;
println!("inner counter: {}", counter);
}
});
let mut counter = 0;
loop {
sleep(Duration::from_millis(1000));
*flag.lock().unwrap() = true;
counter += 1;
if counter > 3 {
break;
}
println!("outside counter: {}", counter);
cond.notify_one();
}
hdl.join().unwrap();
println!("{:?}", flag);
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
输出:
outside counter: 1
inner counter: 1
outside counter: 2
inner counter: 2
outside counter: 3
inner counter: 3
Mutex { data: true, poisoned: false, .. }2
3
4
5
6
7
3.2.6 信号量Semaphore
这里使用的是tokio的Semaohore实现,tokio创建的异步任务。
信号量即对允许使用的任务数量做限制,使用前需要申请信号量,如果容量满了,就需要等待; 使用后需要释放信号量,以便其它等待者可以继续。
use std::sync::Arc;
use tokio::sync::Semaphore;
#[tokio::main]
async fn main() {
let semaphore = Arc::new(Semaphore::new(3));
let mut join_handles = Vec::new();
for _ in 0..5 {
let permit = semaphore.clone().acquire_owned().await.unwrap();
join_handles.push(tokio::spawn(async move {
//
// 在这里执行任务...
//
drop(permit);
}));
}
for handle in join_handles {
handle.await.unwrap();
}
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
3.3 原子操作Atomic与内存顺序
原子指的是一系列不可被 CPU 上下文交换的机器指令,这些指令组合在一起就形成了原子操作。 在多核 CPU 下,当某个 CPU 核心开始运行原子操作时,会先暂停其它 CPU 内核对内存的操作,以保证原子操作不会被其它 CPU 内核所干扰。
由于原子操作是由指令保证的,往往性能会比Mutex等好。
原子操作是无锁的,但是照样会遇到冲突并等待(如CAS下的缓存竞争)。
3.3.1 Atomic作为全局变量
use std::ops::Sub;
use std::sync::atomic::{AtomicU64, Ordering};
use std::thread::{self, JoinHandle};
use std::time::Instant;
const N_TIMES: u64 = 10000000;
const N_THREADS: usize = 10;
static R: AtomicU64 = AtomicU64::new(0);
//原子操作下的u64允许不同线程对其做运算
fn add_n_times(n: u64) -> JoinHandle<()> {
thread::spawn(move || {
for _ in 0..n {
R.fetch_add(1, Ordering::Relaxed);
}
})
}
//fetch_add对其做原子加法,读-加-写是原子操作,此处用到了内存顺序,设置为Relaxed
fn main() {
let s = Instant::now();
let mut threads = Vec::with_capacity(N_THREADS);
for _ in 0..N_THREADS {
threads.push(add_n_times(N_TIMES));
}
for thread in threads {
thread.join().unwrap();
}
assert_eq!(N_TIMES * N_THREADS as u64, R.load(Ordering::Relaxed));
println!("{:?}",Instant::now().sub(s));
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
Atomic的值具有内部不变性,因此无需声明mut。
3.3.2 内存顺序
内存顺序是指 CPU 在访问内存时的顺序,该顺序可能受以下因素的影响:
- 代码中的先后顺序
- 编译器优化导致在编译阶段发生改变(内存重排序 reordering)
- 运行阶段因 CPU 的缓存机制导致顺序被打乱
3.2.2.1 编译器优化导致内存顺序改变
fn main() {
let mut a = 0;
let mut b = 0;
a = 10;
b = 20;
println!("{} {}", a, b);
}2
3
4
5
6
7
8
9
这里编译器翻译的机器指令有可能b在a前面,也有可能不写入内存,直接传入常量
use std::thread;
static mut X: i32 = 0;
fn main() {
let handle = thread::spawn(|| {
unsafe {
let x = X; // 读取并复制 i32 的值
println!("{}", x);
}
});
unsafe {
X = 1;
X = 2;
}
handle.join().unwrap();
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
这里存在潜在的数据竞争风险,编译器可能会合并X的赋值,实际运行输出值为2
3.2.2.2 CPU缓存导致内存顺序改变
CPU 缓存并不一定真正改变指令的执行顺序,而是可能导致不同线程观察到内存读写操作的顺序,与源代码中的顺序不一致。
粗略来说,和CPU的缓存机制有关系。
3.2.2.3 Ordering限制内存顺序的几种规则
枚举成员有:
- Relaxed, 这是最宽松的规则,它对编译器和 CPU 不做任何限制,可以乱序
- Release 释放,设定内存屏障(Memory barrier),保证它之前的操作永远在它之前,但是它后面的操作可能被重排到它前面
- Acquire 获取, 设定内存屏障,保证在它之后的访问永远在它之后,但是它之前的操作却有可能被重排到它后面,往往和Release在不同线程中联合使用
- AcqRel, 是 Acquire 和 Release 的结合,同时拥有它们俩提供的保证。比如你要对一个 atomic 自增 1,同时希望该操作之前和之后的读取或写入操作不会被重新排序
- SeqCst 顺序一致性, SeqCst就像是AcqRel的加强版,它不管原子操作是属于读取还是写入的操作,只要某个线程有用到SeqCst的原子操作,线程中该SeqCst操作前的数据操作绝对不会被重新排在该SeqCst操作之后,且该SeqCst操作后的数据操作也绝对不会被重新排在SeqCst操作前。
这些规则很大程度上是系统提供的,不是编程语言的特性。
不知道选择什么顺序时可以优先选择SeqCst。
3.2.2.4 内存屏障的例子
use std::thread::{self, JoinHandle};
use std::sync::atomic::{Ordering, AtomicBool};
static mut DATA: u64 = 0;
static READY: AtomicBool = AtomicBool::new(false);
fn reset() {
unsafe {
DATA = 0;
}
READY.store(false, Ordering::Relaxed);
}
fn producer() -> JoinHandle<()> {
thread::spawn(move || {
unsafe {
DATA = 100; // A
}
READY.store(true, Ordering::Release); // B: 内存屏障 ↑
})
}
fn consumer() -> JoinHandle<()> {
thread::spawn(move || {
while !READY.load(Ordering::Acquire) {} // C: 内存屏障 ↓
assert_eq!(100, unsafe { DATA }); // D
})
}
fn main() {
loop {
reset();
let t_producer = producer();
let t_consumer = consumer();
t_producer.join().unwrap();
t_consumer.join().unwrap();
}
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
通过Release保证数据写入操作在READY.store之前,通过Acquire保证数据读取在READY.load之后。这里通过while循环确保READY.load在READY.store之后。整体地确保了多线程中最后的读取操作在写入操作之后。
当消费者的 C 读取到生产者 B 写入的 true 时:
其中:
- sb:sequenced-before,线程内的程序顺序。
- sw:synchronizes-with,跨线程同步关系。 最终得到: [ \boxed{A \xrightarrow{hb} D} ]
hb 表示 happens-before。
3.3.3 多线程使用Atomic
配合Arc使用:
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::{hint, thread};
fn main() {
let spinlock = Arc::new(AtomicUsize::new(1));
let spinlock_clone = Arc::clone(&spinlock);
let thread = thread::spawn(move|| {
spinlock_clone.store(0, Ordering::SeqCst);
});
// 等待其它线程释放锁
while spinlock.load(Ordering::SeqCst) != 0 {
hint::spin_loop();
}
if let Err(panic) = thread.join() {
println!("Thread had an error: {:?}", panic);
}
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
3.3.3 原子操作相比锁的不足
那么原子类型既然这么全能,它可以替代锁吗?答案是不行:
对于复杂的场景下,锁的使用简单粗暴,不容易有坑。
std::sync::atomic包中仅提供了数值类型的原子操作:AtomicBool,AtomicIsize,AtomicUsize,AtomicI8,AtomicU16等,而锁可以应用于各种类型。在有些情况下,必须使用锁来配合,例如上一章节中使用
Mutex配合Condvar。
3.3.4 应用场景
事实上,Atomic虽然对于用户不太常用,但是对于高性能库的开发者、标准库开发者都非常常用,它是并发原语的基石,除此之外,还有一些场景适用:
- 无锁(lock free)数据结构
- 全局变量,例如全局自增 ID, 在后续章节会介绍
- 跨线程计数器,例如可以用于统计指标
3.4 基于Send和Sync的线程安全
Rc,RefCell和裸指针无法直接用于多线程。
Rc用于多线程时会报错提示未实现Send特征(尝试将所有权转移至子线程时)。
对比Rc和Arc的源码:
// Rc源码片段
impl<T: ?Sized> !marker::Send for Rc<T> {}
impl<T: ?Sized> !marker::Sync for Rc<T> {}
// Arc源码片段
unsafe impl<T: ?Sized + Sync + Send> Send for Arc<T> {}
unsafe impl<T: ?Sized + Sync + Send> Sync for Arc<T> {}2
3
4
5
6
7
其中!表示移除相应特征实现。
3.4.1 Send和Sync
Send和Sync是 Rust 安全并发的重中之重,但是实际上它们只是标记特征(marker trait,该特征未定义任何行为,因此非常适合用于标记):
- 实现
Send的类型可以在线程间安全的传递其所有权 - 实现
Sync的类型可以在线程间安全的共享(通过引用)
这里还有一个潜在的依赖:如果引用可以安全地在线程间传递,那么其数据就实现了安全的共享(通过引用),即:
此外,还有:
本质是独占访问权(或者说可变访问)的转移,
本质是共享访问权的等价。
Rwlock的实现:
unsafe impl<T: ?Sized + Send + Sync> Sync for RwLock<T> {}由于可以并发读,所以这里的T实现了Sync特征
//Mutex的实现,T没有实现Sync特征
unsafe impl<T: ?Sized + Send> Sync for Mutex<T> {}2
3.4.2 类型的实现
Send和Sync是自动实现类型,对于复合类型,一般子类型实现那么复合类型也自动实现。
| 类型 | Send | Sync | 原因 |
|---|---|---|---|
i32、bool 等基本类型 | ✅ | ✅ | 可以安全转移和共享 |
String | ✅ | ✅ | 所有权转移安全,共享读取安全 |
Rc<T> | ❌ | ❌ | 引用计数不是线程安全的 |
Cell<T> | ✅* | ❌ | 内部可变性缺乏线程同步机制 |
RefCell<T> | ✅* | ❌ | 借用检查不是线程安全的 |
Arc<T> | ✅* | ✅* | 使用原子引用计数,但仍依赖 T 的线程安全性 |
Mutex<T> | ✅* | ✅* | 通过互斥锁保护内部数据 |
*const T、*mut T | ❌ | ❌ | 裸指针不提供线程安全保证 |
注: * 表示有条件成立:
Cell<T>、RefCell<T>:实现Send要求T: Send。Arc<T>:实现Send和Sync都要求T: Send + Sync。Mutex<T>:实现Send和Sync都要求T: Send。
3.4.3 裸指针
裸指针中的地址是虚拟进程空间的地址,一般不是实际的硬件地址。
Rust中裸指针的类型分为*const和*mut,这里的*是类型表达式的一部分,和*ptr中的意义不用。mut不代表一定有效或具有独占访问权。
3.4.3.1 实现Send特征
use std::thread;
#[derive(Debug)]
struct MyBox(*mut u8);
unsafe impl Send for MyBox {}
fn main() {
let p = MyBox(5 as *mut u8);
let t = thread::spawn(move || {
println!("{:?}",p);
});
t.join().unwrap();
}2
3
4
5
6
7
8
9
10
11
12
13
3.4.3.2 实现Sync特征
use std::thread;
use std::sync::Arc;
use std::sync::Mutex;
#[derive(Debug)]
struct MyBox(*const u8);
unsafe impl Send for MyBox {}
unsafe impl Sync for MyBox {}
fn main() {
let b = &MyBox(5 as *const u8);
let v = Arc::new(Mutex::new(b));
let t = thread::spawn(move || {
let _v1 = v.lock().unwrap();
});
t.join().unwrap();
}2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
这里sync使v:Arc<Mutex<&MyBox>>可以在线程间send。
以上例子通过newtype为裸指针实现了Send和Sync特征。