跳到主要内容

Rust 高并发实战(1):所有权、Send/Sync 与并发成本模型

Rainy
雨落无声,代码成诗 —— 致力于技术与艺术的极致平衡
Rainy
8 MIN READ... VIEWS

Rust 的类型系统可以阻止不安全共享,但不会替你决定该共享多少、排队多久和何时拒绝。

本期先区分线程并行和 async I/O,并理解所有权、SendSyncArc 真正保证什么。

架构图:类型安全、调度模型与资源上限是三层问题

Rust 类型系统解决“这份数据能否安全跨线程”,executor 解决“任务何时得到执行”,semaphore/线程池解决“允许多少任务同时占用资源”。三个层次不能互相替代:Send + Sync 的 client 仍可能把数据库连接打满。

一、线程与 async 解决不同等待

工作合适模型上限来源
CPU 密集计算固定线程池/RayonCPU 核数、cache、内存带宽
大量网络 I/OTokio tasksocket、连接池、下游容量、内存
同步阻塞库spawn_blocking/专用线程blocking pool 与独立许可
长生命周期同步循环专用线程 + 有界 channelworker 和队列预算

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>> 不可安全跨线程,编译器会拒绝;
  • 能不共享就不共享,能转移所有权就不加锁。

三、SendSync

  • T: SendT 的所有权可安全移动到另一个线程;
  • 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 结构化并发

参考资料

Logo
RainLib

探索技术、设计与分布式系统的边界。构建面向未来的开发者工具。

留言与建议

© 2026 RainLib. 为未来构建。(Built for the Future)
版权所有。
系统正常