Rust 高并发实战(1):所有权、Send/Sync 与并发成本模型
Rust 的类型系统可以阻止不安全共享,但不会替你决定该共享多少、排队多久和何时拒绝。
本期先区分线程并行和 async I/O,并理解所有权、Send、Sync 与 Arc 真正保证什么。
架构图:类型安全、调度模型与资源上限是三层问题
Rust 类型系统解决“这份数据能否安全跨线程”,executor 解决“任务何时得到执行”,semaphore/线程池解决“允许多少任务同时占用资源”。三个层次不能互相替代:Send + Sync 的 client 仍可能把数据库连接打满。
一、线程与 async 解决不同等待
| 工作 | 合适模型 | 上限来源 |
|---|---|---|
| CPU 密集计算 | 固定线程池/Rayon | CPU 核数、cache、内存带宽 |
| 大量网络 I/O | Tokio task | socket、连接池、下游容量、内存 |
| 同步阻塞库 | spawn_blocking/专用线程 | blocking pool 与独立许可 |
| 长生命周期同步循环 | 专用线程 + 有界 channel | worker 和队列预算 |
async task 在 .await 未就绪时保存状态并让出 worker,因此等待很经济;但每个 task 仍保留 Future 状态、参数、buffer、timer 和外部连接。
二、所有权把共享变成显式设计
use std::sync::{Arc, Mutex};
let state = Arc::new(Mutex::new(State::default()));
let worker_state = Arc::clone(&state);
std::thread::spawn(move || {
worker_state.lock().unwrap().completed += 1;
});
Arc<T>提供线程安全的共享所有权,不自动让T内部线程安全;Mutex<T>提供互斥和内部可变性;Rc<RefCell<T>>不可安全跨线程,编译器会拒绝;- 能不共享就不共享,能转移所有权就不加锁。
三、Send 与 Sync
T: Send:T的所有权可安全移动到另一个线程;T: Sync:&T可安全在线程间共享;- 多线程 Tokio 的
spawn通常要求 Future 为Send + 'static,因为 task 可能在另一个 worker 上继续 poll; - 手工
unsafe impl Send/Sync等于向编译器承诺底层不变量,必须有严谨证明和测试。
async fn bad() {
let value = std::rc::Rc::new(1);
tokio::task::yield_now().await;
println!("{value}");
}
// tokio::spawn(bad()) 在多线程 runtime 中不会通过 Send 约束。
这不是 async “不能用 Rc”,而是这个 Future 可能跨线程。单线程 LocalSet 可运行 !Send task,但要明确它会共享单个执行线程的容量与故障域。
四、safe Rust 仍可能写出的并发错误
- 两把 Mutex 逆序获取导致死锁;
- 无界 channel 或无界
tokio::spawn导致 RSS 失控; - 同步阻塞占住 runtime worker;
- timeout 后 detached task 继续产生副作用;
- atomic ordering 不满足发布协议;
- 业务 read-modify-write 跨步骤失去原子性。
内存安全只是并发正确性的子集。
五、容量模型
入口 10,000 RPS、每请求 3 个下游、平均 20 ms:
30,000 calls/s × 0.020 s = 600 个平均在途下游调用
将 CPU、内存、连接和 deadline 分别预算。下游 P99 200 ms 时不能简单把 semaphore 扩到 6,000;数据库只有 200 连接时,额外 task 只会转移到连接池排队。
每个实例的许可还要乘以副本数审计总量。滚动发布期间副本暂时增加,也可能让总连接超过数据库上限。
六、线程实验
use std::{sync::Arc, thread};
use std::sync::atomic::{AtomicU64, Ordering};
fn main() {
let total = Arc::new(AtomicU64::new(0));
let mut workers = Vec::new();
for _ in 0..4 {
let total = total.clone();
workers.push(thread::spawn(move || {
for _ in 0..1_000_000 {
total.fetch_add(1, Ordering::Relaxed);
}
}));
}
for worker in workers { worker.join().unwrap(); }
assert_eq!(total.load(Ordering::Relaxed), 4_000_000);
}
Relaxed 足以保证单个计数不丢更新,但不为其他数据建立发布顺序。后续会用 Loom 验证 Acquire/Release 协议。
七、Scoped Thread 与所有权边界
普通 thread::spawn 要求 closure 为 'static,因为线程可能活得比当前栈帧更久。std::thread::scope 让线程在 scope 结束前必定 join,因此可以安全借用局部数据:
let mut chunks = vec![0_u64; 8];
std::thread::scope(|scope| {
for chunk in &mut chunks {
scope.spawn(move || {
*chunk = expensive_calculation();
});
}
});
assert!(chunks.iter().all(|value| *value > 0));
这体现了结构化并发:子线程不能逃出父作用域。对 CPU 批处理,它比 Arc<Mutex<Vec<_>>> 更少共享、更容易证明。
八、Arc 的成本与引用环
Arc::clone/drop 需要 atomic 引用计数,会产生 cache line 通信。它通常不是首要瓶颈,但在每消息多次 clone 的热路径中应 profile。把 Arc clone 移到请求外层、传 &T 或批量处理可减少成本。
Arc 不会自动打破环:A 强引用 B、B 强引用 A 会永久保留。父子、订阅者和 callback 图中使用 Weak 表达非拥有关系,并测试连接关闭后对象计数回落。
九、False Sharing 与数据布局
两个线程频繁写不同 atomic,若它们落在同一 cache line,仍会反复争夺一致性所有权。症状是 CPU 高、扩核后吞吐不升、perf 显示原子/缓存相关开销。
优化顺序:先减少共享写,再做每线程局部计数并周期聚合,最后才考虑 cache padding。padding 依赖硬件布局,不能替代架构优化。
十、Rust 内存模型的工程边界
safe Rust 防止未同步非原子访问造成内存 data race;atomic ordering 仍要由算法证明。Relaxed 仅保证该原子操作本身不可撕裂与修改顺序,不发布旁边普通数据。Acquire/Release 建立跨线程可见性,SeqCst 额外提供全局单序,但“更强”不能修复错误状态机。
涉及 unsafe、FFI、DMA 或自定义并发结构时,还要审计:别名规则、对象生命周期、回调线程、C 库线程安全和 panic 穿越边界。
十一、选择线程、Task 还是 Actor
工作是否等待 async I/O? ─是→ Tokio Future
│否
↓
是否持续占用 CPU? ─是→ Rayon/固定 CPU pool
│否
↓
是否调用短期同步阻塞库? ─是→ spawn_blocking + gate
│否
↓
是否长期持有同步资源? ─是→ 专用 thread + bounded channel
Actor 不是第四种调度器,而是状态所有权模式;它可以运行在 Tokio task 或专用线程。选择依据是状态是否适合单 owner 顺序处理,以及 mailbox 能否有界。
十二、容量表必须包含 task 大小
Rust Future 大小取决于跨 await 保存的字段。用 std::mem::size_of_val(&future) 可做局部实验,但真实内存还包括 task header、Arc、allocator 和外部 buffer。用进程 RSS/heap profiler在稳态和峰值验证,不能只乘 size_of。
12.1 用所有权减少共享,而不是用 Arc 包住一切
批处理可以把输入拆成互不重叠的 chunk,把所有权移动给 worker,完成后再聚合。只有配置、client pool 等真正共享的长寿命对象需要 Arc。请求局部结果通过返回值传递,通常比 Arc<Mutex<Vec<Result>>> 更简单。
let handles: Vec<_> = inputs
.chunks(batch_size)
.map(|chunk| {
let owned = chunk.to_vec();
std::thread::spawn(move || process_batch(owned))
})
.collect();
let results: Result<Vec<_>, _> = handles
.into_iter()
.map(|handle| handle.join().map_err(|_| Error::WorkerPanic))
.collect();
这里每个 worker 独占自己的 Vec,不存在共享写入,也不需要锁。内存代价是复制;如果复制昂贵,使用 scoped thread 借用不重叠 slice。
12.2 证明边界
为共享结构写明不变量:谁能修改、修改需要什么 guard、何时发布给读者、panic/cancel 后是否仍成立。unsafe impl Send/Sync 必须把这些证明变成 Safety 注释和模型测试,否则只是关闭编译器告警。
十三、本期验收
- 把服务任务分成 async I/O、CPU、同步阻塞和长期线程;
- 计算目标 RPS 下平均/P99 在途数和内存/连接预算;
- 写一个
Rc跨线程失败示例并解释编译器错误; - 列出 safe Rust 仍需测试的业务不变量。
下一期:Rust-2:Future、Poll/Waker 与 Tokio 结构化并发。