跳到主要内容

Rust 高并发实战(4):锁、Channel、Atomic 与取消安全

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

编译通过证明内存访问安全,不证明业务事务完整、锁不会死锁或 Future 可以在任意 await 点被取消。

原理图:共享状态、消息所有权与外部事务

Mutex 适合短内存不变量,actor 适合复杂单 owner 状态,但跨进程副作用最终仍需数据库事务、幂等键或 outbox。把状态放进 actor 并不会让网络写入自动原子。

一、同步原语选择表

场景首选注意
极短临界区、不跨 awaitstd::sync::Mutex竞争时会阻塞 worker
guard 必须跨 awaittokio::sync::Mutex尽量重新设计以缩短持锁
单资源复杂状态actor + mpscchannel 必须有界
高频读、整体更新不可变 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:阻塞隔离与性能诊断

参考资料

Logo
RainLib

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

留言与建议

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