重构与思考:工业级基础库中的协作式取消 (Cooperative Cancellation)
在构建高吞吐、低延迟的分布式系统时,如何优雅地终止正在运行的复杂任务往往比启动它更具挑战。如果任务涉及网络 I/O、磁盘读写或复杂计算,简单粗暴地 kill 线程不仅导致资源泄露——文件句柄未关闭、锁未释放——还会破坏数据一致性。
某工业级分布式基础库通过基于 Future/Promise 机制的"协作式取消令牌(Cooperative Cancellation Token)"设计解决此问题。这一模式值得深入理解其设计哲学与权衡。
什么是协作式取消?
"协作式(Cooperative)"意味着任务自身主动检查并响应终止信号,而非被外部强制关闭。
想象一场会议:如果老板直接关灯(强制终止),大家陷入混乱——笔记本没合,水杯打翻。协作式做法是,老板看一眼手表,给一个眼神,大家心领神会,收拾东西,有序离场。
代码中涉及两个角色:
- 发起方 (Source):持有"开关",决定何时发出取消信号。
- 执行方 (Token):持有令牌,在关键执行点(检查点)检查其状态。
深度解析:基于 Future 的信号传递
此库没有为取消逻辑单独发明锁或条件变量,而是复用现有异步基础设施。它利用异步框架中的 Promise<void> 和 Future<void>。
机制:
- Source (CancellationTokenSource):持有
Promise<void>。调用Cancel()时,通过Promise::SetValue()设置承诺。 - Token (CancellationToken):持有对应的
Future<void>。 - Check (IsCancellationRequested):测试 future 是否已就绪。
权衡分析:
优势:
- 语义统一:取消成为标准异步事件。等待取消信号看起来像等待任何其他异步结果。
- 自然组合:
WaitAny和WhenAll等 future 组合算子无需额外代码就能处理取消。你写出WaitAny(NetworkFuture, CancellationFuture)来表达"要么网络完成,要么任务被取消"——不需要自定义轮询逻辑。
劣势:
- 资源成本:每个令牌需要一个共享状态块。拥有数百万短期任务的系统中,内存开销成为问题。
- 层级复杂性:从深层任务树派生子令牌(链接令牌)引入挑战。
净室重构:Rust 视角下的复述
为了更清晰地演示"source-token 分离"与共享状态模式,我们用 Rust 重构设计,剥离 Future 包装以暴露核心逻辑。
注:这是模式的演示,非生产代码。
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::Duration;
/// 协作式取消的核心:状态共享
/// Source 持有写入权,Token 持有读取权
pub struct MyCancellationTokenSource {
shared: Arc<AtomicBool>,
}
impl MyCancellationTokenSource {
pub fn new() -> Self {
Self {
shared: Arc::new(AtomicBool::new(false)),
}
}
/// 派发一个只读的令牌给任务方
pub fn token(&self) -> MyCancellationToken {
MyCancellationToken {
shared: self.shared.clone(),
}
}
/// 发起方:按下停止按钮
pub fn cancel(&self) {
self.shared.store(true, Ordering::SeqCst);
}
}
pub struct MyCancellationToken {
shared: Arc<AtomicBool>, // 共享的原子布尔值
}
impl MyCancellationToken {
/// 任务方:非阻塞检查
pub fn is_cancellation_requested(&self) -> bool {
self.shared.load(Ordering::SeqCst)
}
/// 任务方:模拟“如果取消则抛出异常/错误”的语义
pub fn check(&self) -> Result<(), String> {
if self.is_cancellation_requested() {
Err("Operation cancelled".to_string())
} else {
Ok(())
}
}
}
fn main() {
let source = MyCancellationTokenSource::new();
let token = source.token();
println!("[Main] Starting worker thread...");
let handle = thread::spawn(move || {
for i in 0..10 {
// 关键点:协作式检查
// 任务必须在合适的时机主动询问“我还需要继续吗?”
if let Err(e) = token.check() {
println!("[Worker] Detected cancellation: {}", e);
return;
}
println!("[Worker] Processing step {}...", i);
thread::sleep(Duration::from_millis(200));
}
println!("[Worker] Task completed successfully.");
});
// 模拟运行一段时间后取消
thread::sleep(Duration::from_millis(700));
println!("[Main] Requesting cancellation...");
source.cancel();
handle.join().unwrap();
println!("[Main] Program exited.");
}
代码解读
所有权分离:
MyCancellationTokenSource负责令牌创建和状态修改;MyCancellationToken仅读状态。这遵循单一职责原则,防止执行方意外修改取消状态。原子性:
AtomicBool配合Ordering::SeqCst保证多线程可见性。工业级实现通常通过内存屏障或更轻量级的排序(如在同步点使用Relaxed)优化。检查语义:
check()方法模拟原库的ThrowIfCancellationRequested()。Rust 用Result替代异常,与显式错误处理哲学一致。
核心模式:通过共享状态的单向信号流。 理解 source 如何设置信号、token 如何读取信号,你就能推理取消在系统中的代价,以及它在任务架构中的位置。