Rust 高并发实战(4):锁、Channel、Atomic 与取消安全
编译通过证明内存访问安全,不证明业务事务完整、锁不会死锁或 Future 可以在任意 await 点被取消。
原理图:共享状态、消息所有权与外部事务
Mutex 适合短内存不变量,actor 适合复杂单 owner 状态,但跨进程副作用最终仍需数据库事务、幂等键或 outbox。把状态放进 actor 并不会让网络写入自动原子。
一、同步原语选择表
| 场景 | 首选 | 注意 |
|---|---|---|
| 极短临界区、不跨 await | std::sync::Mutex | 竞争时会阻塞 worker |
| guard 必须跨 await | tokio::sync::Mutex | 尽量重新设计以缩短持锁 |
| 单资源复杂状态 | actor + mpsc | channel 必须有界 |
| 高频读、整体更新 | 不可变 Arc 快照 | 发布后不可内部修改 |
| 独立计数/标志 | Atomic | 多字段不变量不是事务 |
Tokio 官方建议不要在 async 代码里条件反射地使用 async mutex;短且低竞争的临界区,同步 Mutex 可以接受。红线是持有同步 guard 跨 .await。
二、锁外 I/O 与重复加载
let cached = state.cache.lock().unwrap().get(&key).cloned();
if let Some(value) = cached { return Ok(value); }
let loaded = fetch_remote(&key).await?;
let mut cache = state.cache.lock().unwrap();
Ok(cache.entry(key).or_insert(loaded).clone())
锁外 I/O 避免全局等待,但两个 task 可能重复 fetch。按业务使用 per-key singleflight、actor 或允许重复。不要为了消除重复把远程 I/O 搬回锁内。
三、Actor:让状态只有一个 owner
enum Command {
Get { key: String, reply: oneshot::Sender<Option<Value>> },
Put { key: String, value: Value, reply: oneshot::Sender<()> },
}
async fn state_actor(mut rx: mpsc::Receiver<Command>) {
let mut data = HashMap::new();
while let Some(command) = rx.recv().await {
match command {
Command::Get { key, reply } => { let _ = reply.send(data.get(&key).cloned()); }
Command::Put { key, value, reply } => { data.insert(key, value); let _ = reply.send(()); }
}
}
}
actor 消除共享锁,但会形成串行瓶颈。监控 mailbox 深度/等待;需要时按 key 分片多个 actor。
四、取消安全与副作用
async fn transfer(from: Id, to: Id, amount: Money) -> Result<()> {
debit(from, amount).await?;
credit(to, amount).await?; // 此处取消会留下半笔操作
Ok(())
}
timeout/select drop Future 不会回滚已提交副作用。用数据库事务、幂等键、outbox 或状态机:
Pending → Debited → Credited → Completed
↘ Compensating → Compensated
逐个检查 channel send/recv、流读取和锁获取的 cancellation safety 文档,并测试每个 await 前后取消。
五、Atomic 与 Loom
#[test]
fn publication_is_visible() {
loom::model(|| {
let ready = Arc::new(AtomicBool::new(false));
let value = Arc::new(AtomicUsize::new(0));
let r1 = ready.clone(); let v1 = value.clone();
let writer = thread::spawn(move || {
v1.store(42, Ordering::Relaxed);
r1.store(true, Ordering::Release);
});
let r2 = ready.clone(); let v2 = value.clone();
let reader = thread::spawn(move || {
if r2.load(Ordering::Acquire) {
assert_eq!(v2.load(Ordering::Relaxed), 42);
}
});
writer.join().unwrap(); reader.join().unwrap();
});
}
Loom 枚举小型并发协议的可能交错。只模型化核心结构,控制分支和线程数;把 ordering 证明写在封装 atomic 的模块中。
六、Mutex Poisoning 与错误策略
std::sync::Mutex 在持锁线程 panic 后会 poisoned,提示受保护状态可能只更新了一半。直接 unwrap() 会把故障传播成更多 panic;盲目 into_inner() 又可能继续使用破坏的不变量。
按状态性质选择:缓存可清空重建;账务状态应停止服务并从事务源恢复;只更新独立计数可能验证后继续。把 poison recovery 写成显式函数和告警,而不是散落的 unwrap。
Tokio Mutex 不采用同样的 poisoning 语义,因此 panic 时 guard drop 会释放锁,但业务状态仍可能处于中间步骤。内存可访问不代表不变量正确。
七、死锁之外的活锁与饥饿
- 死锁:互相等待,永不前进;
- 活锁:不断重试/让步但没有完成;
- 饥饿:某些任务长期拿不到资源;
- 优先级反转:高优任务等待低优任务持有的资源。
try_lock 紧循环会制造活锁和 CPU 空转。失败后要有有界退避/jitter,或改为正常 await。公平队列改善饥饿,但会带来 head-of-line blocking,仍需缩短临界区。
八、Channel 协议不只是类型
为每条 channel 定义:谁发送、谁关闭、容量、满时策略、消息是否可丢、receiver 重启后的语义。oneshot 适合一次回复;mpsc 适合多生产者单 owner;broadcast 的慢订阅者会 lag,需要处理丢消息;watch 表达“只关心最新状态”。
关闭 sender 只是没有新消息,receiver 仍应 drain 已排队消息;服务停机是否 drain 取决于消息是否可重放和 shutdown deadline。
九、Loom 模型设计技巧
状态空间会指数增长:限制到 2~3 线程、少量循环和最小状态机。不要在模型中使用标准库 Mutex/Atomic,必须使用 loom 对应类型才能控制调度。
一个好的 Loom 测试断言业务不变量,而不是某个预期执行顺序,例如“读到 ready 后 value 必为 42”“最多一个线程从 Initial 成功转为 Leader”。失败时保存交错并变成长期回归。
十、并发 API 评审清单
- 类型是否真的需要
Arc<Mutex<T>>; - guard 生命周期是否跨 await/回调;
- 多锁顺序是否唯一;
- channel 是否有界并处理关闭;
- atomic ordering 是否有注释与模型测试;
- panic/cancel 发生在中间步骤时状态是否可恢复;
- unsafe/FFI 是否维持 Send/Sync 与生命周期承诺。
十一、分片状态与一致性边界
DashMap/分片 Mutex 可以降低单 key 操作争用,但跨 key 原子操作仍需更高层协议。转账若账户落在不同 shard,同时持两把锁必须按稳定 key 排序;更推荐把强一致不变量下沉到数据库事务。
热点 key 不会因 64 分片自动消失。对热点使用 singleflight、复制只读数据、请求合并或业务分区,并监控每 shard 等待而不只看总体平均。
十二、Panic 与 Task 错误
spawn task 的 panic 包装在 JoinError;丢弃 JoinHandle 会丢失观测。服务边界决定 panic 策略:请求 task panic 可返回 500 并报警,核心 actor panic 可能需要停止 readiness 和重建状态。不要用 catch_unwind 掩盖已经破坏的不变量。
12.1 取消安全状态机
每个持久状态都必须可重入:重复执行 Debited → Completed 不会重复入账;进程崩溃后 worker 能扫描 Debited 并继续或补偿。Future 是否还在内存中不再是业务正确性的唯一依据。
12.2 两锁顺序的证明
账户转账需要同时锁两账户时,以稳定 ID 排序获取:先 min(from,to),再 max(from,to)。所有代码路径都使用同一 helper,不能靠调用方记忆。Loom 测试同时运行 A→B 与 B→A,断言二者最终完成且余额守恒。
12.3 Async Mutex 持锁跨 await 的成本
即使 Tokio Mutex 允许 guard 跨 await,远程 500 ms 长尾也会让后续 waiter 全部排队。将远程读移出锁后,再以版本/CAS 提交;若必须串行远程操作,actor mailbox 的队列深度和超时比隐式锁等待更易治理。
十三、本期验收
- 比较 Mutex、async Mutex、actor 在真实读写比下的吞吐和 P99;
- 故意制造两把锁逆序死锁并建立统一锁顺序;
- 对转账流程在每个 await 点注入取消;
- 用 Loom 验证一个发布或 CAS 状态机。
上一篇:Rust-3:Axum 有界服务 · 下一篇:Rust-5:阻塞隔离与性能诊断。