Rust 异步 Cancel Safety 与优雅停机工程实战:从跨 .await 灾难到生产级优雅关闭
每一个跑过 async Rust 的人都被半夜的 oncall 叫醒过——服务重启时数据库事务只写了一半,消息队列里出现了幽灵消息,下游服务收到了残缺的请求。元凶不是 borrow checker 不够聪明,是我们忽视了 Cancel Safety 这个藏在
.await背后的暗坑。
一、Cancel Safety:async 世界最被忽视的语义陷阱
Rust 的 async/await 模型在语法层面极其优雅,但过早取消(early cancellation)的语义却与同步代码截然不同。在同步代码中,一段代码要么执行完毕,要么根本不执行;而在 async 代码中,Future 可能在任意 .await 点被 Drop 掉,导致后续代码永远不执行。
这不是 bug,而是设计选择——但问题在于,标准库和很多第三方库的 API 注释里从来不写清楚某个 Future 是否是 Cancel Safe 的。
1.1 一个看似无害的 .await
async fn process_message(db: &DbPool, msg: Message) -> Result<()> {
let mut txn = db.begin_txn().await?; // await 点 A
txn.insert(&msg).await?; // await 点 B
txn.commit().await?; // await 点 C
Ok(())
}
当 tokio 运行时收到 SIGTERM,它开始 Drop 未完成的 task。假设 process_message 正在 .await point B 处等待 I/O,此时 Drop 会被调用——后果是:事务被放弃,连接归还到池里(可能持有未提交的更改),消息"丢失"了。
1.2 Cancel Safe 的严格定义
一个 Future 是 Cancel Safe 的,当且仅当:
- 从它被首次 Poll 开始,到它被 Drop 为止,它所执行的操作要么完全不生效,要么完全生效——不存在中间状态;
- Drop 它不会导致资源泄漏、数据不一致或死锁。
标准库中,tokio::io::AsyncReadExt::read_to_end 就是非 Cancel Safe 的典型例子——如果你在它读完之前 Drop,部分数据已经被读走,但返回的却是 Err(Cancelled),你不知道读了多少。
1.3 四种常见的 Cancel 非安全模式
| 模式 | 表现 | 修复策略 |
|---|---|---|
| 操作原子性破坏 | 半完成的数据库事务 | scopeguard 或手动 Drop 实现 |
| 消息重复消费 | MQ ack 未发送,消息被重新投递 | 设计为幂等消费 + at-least-once |
| 文件状态不一致 | 写入部分数据后取消 | O_CREATE|O_EXCL + 原子 rename |
| 死锁 | 持有 MutexGuard 时 Drop | tokio::sync::Mutex + Drop-aware wrapper |
二、Tokio 取消机制的运行时真相
理解 tokio 的取消实现方式是写出正确优雅停机代码的前提。
2.1 tokio::time::timeout 内部做了什么
pub async fn timeout<T>(duration: Duration, future: T) -> Result<T::Output, Elapsed>
where
T: Future,
{
pin!(future);
pin!(sleep(duration));
select! {
_ = &mut sleep => Err(Elapsed),
res = &mut future => Ok(res),
}
}
tokio::select! 和 tokio::time::timeout 在底层都是基于 Future::poll 的"竞争"。当一个分支先 ready,其他分支直接被 Drop。这意味着:被 Drop 的 Future 没有任何机会执行 Drop 后逻辑——它不会收到任何回调,不会有机会清理。
2.2 JoinHandle::abort 的行政命令
let handle = tokio::spawn(async { heavy_computation().await });
// 运维脚本触发:
handle.abort();
// handle 返回 JoinError::cancelled
abort() 不做任何"通知"——它直接对 task 的 Future 执行 Drop。如果你的 task 正在 .await 一个数据库查询,那这次查询会被直接放弃。
关键问题:abort 不是协作式的,它不给 task 留任何退出路径。这就是为什么生产中我们必须使用协作式取消原语。
2.3 CancellationToken:协作式取消的基石
Tokio 官方提供的 tokio_util::sync::CancellationToken 是优雅停机的核心原语:
use tokio_util::sync::CancellationToken;
let token = CancellationToken::new();
let child_token = token.child_token();
// Worker task
let worker = tokio::spawn(async move {
loop {
tokio::select! {
_ = token.cancelled() => {
// 执行清理逻辑
break;
}
conn = listener.accept() => {
// 处理连接
}
}
}
});
// 触发取消
token.cancel();
worker.await.unwrap();
注意 cancelled() 返回的 Future 是 Cancel Safe 的——它只返回 () 或永远 Pending,不存在中间状态。
三、优雅停机三阶段协议
一个生产级服务不是"收到信号就退出",而是要走完三个阶段:
3.1 Phase 1:停止接受新请求
/// 停止接受新连接,但标记 shutdown 状态
async fn enter_shutdown_mode(listener: &TcpListener, cancel_token: &CancellationToken) {
// 关闭监听 socket 文件描述符
// 通知负载均衡器 /health 返回 503
cancel_token.cancelled().await;
// 给 upstream 一个宽限期去摘除 endpoints
tokio::time::sleep(Duration::from_secs(5)).await;
}
3.2 Phase 2:等待 in-flight 请求完成
/// 实时跟踪活跃连接数的 Atomic 计数器
struct GracefulShutdown {
active_connections: Arc<AtomicU32>,
max_wait: Duration,
}
impl GracefulShutdown {
async fn wait_for_drain(&self) -> Result<()> {
let deadline = Instant::now() + self.max_wait;
loop {
let count = self.active_connections.load(Ordering::Relaxed);
if count == 0 {
return Ok(());
}
if Instant::now() >= deadline {
return Err(anyhow!("drain timeout, {} still active", count));
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
}
3.3 Phase 3:释放资源并退出
async fn release_resources(
db_pool: DbPool,
redis: RedisClient,
metrics: MetricsExporter,
) -> Result<()> {
// 1. 刷新 metrics
metrics.flush().await?;
// 2. 关闭连接池(等待归还所有连接)
db_pool.close().await;
// 3. 关闭 Redis 管道
redis.close().await;
Ok(())
}
完整的优雅停机编排:
async fn graceful_shutdown(
handles: Vec<JoinHandle<()>>,
shutdown: GracefulShutdown,
db_pool: DbPool,
redis: RedisClient,
) -> Result<()> {
log::info!("Shutdown initiated, stopping new connections...");
// Phase 2: 等待排空
shutdown.wait_for_drain().await?;
// Phase 3: 释放资源
release_resources(db_pool, redis, MetricsExporter::global()).await?;
// 等待所有 worker 退出(应该很快)
for handle in handles {
if let Err(e) = handle.await {
if !e.is_cancelled() {
log::error!("Worker panicked: {}", e);
}
}
}
log::info!("Shutdown complete");
Ok(())
}
四、信号处理:Unix Signal 与 Task 生命周期的桥梁
类 Unix 系统优雅停机需要处理至少三个信号:SIGTERM(k8s 的默认优雅停机信号)、SIGINT(Ctrl-C)和 SIGUSR1(日志轮转或热重载)。
4.1 一站式 Signal Handler
use tokio::signal::unix::{signal, SignalKind};
async fn wait_for_shutdown_signal() -> &'static str {
let mut sigterm = signal(SignalKind::terminate()).unwrap();
let mut sigint = signal(SignalKind::interrupt()).unwrap();
let mut sigusr1 = signal(SignalKind::user_defined1()).unwrap();
tokio::select! {
_ = sigterm.recv() => "sigterm",
_ = sigint.recv() => "sigint",
_ = sigusr1.recv() => "sigusr1",
}
}
4.2 基于 signal 的 Kahn Process
#[tokio::main]
async fn main() {
let state = Arc::new(AppState::new().await);
let mut worker_handles = vec![];
// 启动工作线程
for i in 0..num_cpus::get() {
let state = Arc::clone(&state);
worker_handles.push(tokio::spawn(async move {
worker_loop(i, state).await;
}));
}
// 等待停机信号
let sig = wait_for_shutdown_signal().await;
match sig {
"sigterm" | "sigint" => {
log::info!("Received {}, starting graceful shutdown", sig);
graceful_shutdown(worker_handles, state.shutdown.clone(),
state.db_pool.clone(), state.redis.clone())
.await
.expect("shutdown failed");
}
"sigusr1" => {
log::info!("Received SIGHUP, rotating logs");
state.log_reopen().await;
}
_ => unreachable!(),
}
}
五、生产级优雅停机的五个暗坑
5.1 暗坑一:Mutex 跨越 Drop 边界
std::sync::MutexGuard 不是 Future-aware 的——如果你在 .await 前获取了它并在 .await 之后 Drop,而 await 时任务被 abort,MutexGuard 永远不会释放,直接死锁。
// 危险代码:MutexGuard 跨 await,可能永远不释放
async fn dangerous(pool: &ConnectionPool) {
let guard = pool.lock().await; // acquire
let conn = pool.get_conn().await; // <- 如果在这被 abort
use_conn(conn).await;
drop(guard); // <- 永远到不了这里
}
修复方案:要么使用 tokio::sync::Mutex(它跨 Drop 安全),要么确保临界区没有 .await 点。
5.2 暗坑二:Drop Guard 的取消屏障
有时你需要"不可取消的临界区"——无论信号如何,都要完成。tokio 原生不支持,但可以用 CancellationToken::run_until_cancelled 反转语义:
/// 在最关键的事务提交阶段禁用取消
async fn atomic_commit(txn: &mut Transaction) -> Result<()> {
// 忽略 CancellationToken,强制完成
CancellationGuard::protect_section(async {
txn.write_final_record().await?;
txn.commit().await?;
Ok(())
})
.await
}
自定义实现:
use std::pin::Pin;
use std::task::{Context, Poll};
struct CancellationGuard<F> {
inner: F,
shielded: bool,
}
impl<F: Future> Future for CancellationGuard<F> {
type Output = F::Output;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
// 临时将 waker 替换为空操作 waker,屏蔽取消
let _guard = noop_wenter(cx);
unsafe { self.map_unchecked_mut(|s| &mut s.inner) }.poll(cx)
}
}
警告:noop waker hack 是 UB-adjacent 的代码。更安全的做法是提升这段代码到 spawn_blocking,或者用
tokio::task::unconstrained(nightly)。
5.3 暗坑三:连接池的"僵尸归还"
// 反模式:async Drop + 连接归还
struct PooledConnection {
pool: Arc<DbPool>,
conn: Option<Connection>,
}
// tokio 不支持 async Drop!以下代码不会工作:
// impl Drop for PooledConnection { async fn drop(...) } // 编译错误
正确的归还方式是 tokio::spawn 一个后台 task 来做归还,或者使用 tokio_util::sync::ReusableBoxFuture 这类结构。
5.4 暗坑四:gRPC streaming 的优雅关闭
gRPC 的 Streaming<Request> 在 shutdown 时需要特殊处理——直接 Drop 会导致 RST_STREAM 而不是更优雅的 GOAWAY:
async fn graceful_grpc_shutdown(
server: Server,
cancel_token: CancellationToken,
) {
// 1. 发送 GOAWAY 告诉客户端不要再开新 stream
server.graceful_shutdown().await;
// 2. 等待 in-flight RPC 完成(有超时兜底)
tokio::select! {
_ = wait_all_rpcs_done() => {},
_ = tokio::time::sleep(Duration::from_secs(30)) => {
log::warn!("Forced shutdown, aborting remaining RPCs");
server.shutdown().await;
}
}
}
5.5 暗坑五:Metrics 导出器的最终 flush
很多团队忘记在 shutdown 时 flush 自定义 metrics,导致最后一分钟的时间窗口数据丢失:
impl Drop for PipelineGuard {
fn drop(&mut self) {
// Drop 中不能 await,所以 spawn 一个 blocking task
if let Some(flush_fn) = self.flush.take() {
let _ = std::thread::Builder::new()
.name("final-metrics-flush".into())
.spawn(flush_fn);
}
}
}
六、一个最小可工作的生产模板
把上面所有内容串成一个完整的入口模板:
// src/main.rs
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
use std::time::Duration;
use tokio::time::timeout;
#[tokio::main]
async fn main() -> Result<()> {
// 初始化基础设施
let db_pool = init_db_pool().await?;
let redis = init_redis().await?;
let active_reqs = Arc::new(AtomicU32::new(0));
let shutdown_token = CancellationToken::new();
// 启动 HTTP server
let server_handle = tokio::spawn({
let token = shutdown_token.clone();
let counter = active_reqs.clone();
async move {
axum_server(token, counter).await;
}
});
// 启动 gRPC server
let grpc_handle = tokio::spawn({
let token = shutdown_token.clone();
async move {
grpc_server(token).await
}
});
// 启动 background workers
let worker_handles: Vec<_> = (0..4)
.map(|id| {
let token = shutdown_token.clone();
tokio::spawn(async move {
background_worker(id, token).await
})
})
.collect();
// 等待退出信号
match wait_for_signal().await {
"sigterm" | "sigint" => {
perform_shutdown(
shutdown_token,
vec![server_handle, grpc_handle],
worker_handles,
active_reqs,
db_pool,
redis,
).await?;
}
_ => {}
}
Ok(())
}
async fn perform_shutdown(
token: CancellationToken,
servers: Vec<JoinHandle<()>>,
workers: Vec<JoinHandle<()>>,
active_reqs: Arc<AtomicU32>,
db_pool: DbPool,
redis: RedisClient,
) -> Result<()> {
const DRAIN_TIMEOUT: Duration = Duration::from_secs(25);
// Phase 1: 取消 Token,停止接受新连接
token.cancel();
tokio::time::sleep(Duration::from_secs(2)).await; // 给 LB 摘除窗口
// Phase 2: 等待排空
let drain_result = timeout(DRAIN_TIMEOUT, async {
while active_reqs.load(Ordering::Relaxed) > 0 {
tokio::time::sleep(Duration::from_millis(100)).await;
}
}).await;
if drain_result.is_err() {
let remaining = active_reqs.load(Ordering::Relaxed);
metrics::counter!("graceful_shutdown.drain_timeout", 1);
log::warn!("Drain timeout, {} requests still in-flight", remaining);
}
// Phase 3: 释放资源
timeout(Duration::from_secs(5), db_pool.close()).await??;
timeout(Duration::from_secs(3), redis.close()).await??;
// 等待 handle 完成
for h in servers.into_iter().chain(workers.into_iter()) {
if let Err(e) = timeout(Duration::from_secs(2), h).await {
log::error!("Handle didn't exit in time, aborting");
}
}
Ok(())
}
七、可观测性:别让优雅停机变成黑盒
优雅停机本身必须可观测,否则凌晨三点的 oncall 永远不知道为什么重启卡住了:
// 关键 metrics
metrics::gauge!("shutdown.active_requests", count);
metrics::counter!("shutdown.drain_timeout", 1);
metrics::histogram!("shutdown.drain_duration_secs", elapsed);
metrics::counter!("shutdown.phase1_start", 1);
metrics::counter!("shutdown.phase2_drain_complete", 1);
metrics::counter!("shutdown.phase3_released", 1);
在 Grafana 中配置 Shutdown Panel,监控四个核心指标:从 SIGTERM 到接受的请求数降为 0 的耗时、超时触发次数、最后释放的资源类型。
八、思路延伸:Cancel Safety 对系统设计的深层影响
最后把 Cancel Safety 的影响从代码层面拉升到架构层面:
1. 尽可能使用幂等设计
Cancel 随时发生意味着 at-least-once 交付。如果你的下游不幂等,最终数据会重复。给出一个全局唯一的 request_id 并在下游去重,是基本要求。
2. 分布式事务的三阶段化 传统两阶段提交在协调者 Cancel 时会阻塞。Saga pattern 天然兼容 Cancel——每个步骤都有对应补偿(compensation)步骤,Cancel 就是触发补偿。
3. 架构风格的自然结论
Cancel 是 async 世界的光速壁垒——你无法在执行过程中"回滚时间"。正确的设计使命是:让 Cancel 只发生在"无副作用的等待点",并通过幂等设计吸收 at-least-once 的冲击。
结语
优雅停机不是流程文档上的 checklist,而是代码路径上的工程承诺。从理解 Cancel Safety 的基础语义,到正确使用 CancellationToken 构建协作式取消,再到编排完整的三阶段停机——每一步都是 async Rust 生产中必须跨过去的坎。
下次再被 oncall 叫醒,希望是喝咖啡拖慢了 flush 的速度,而不是凌晨三点的 P0。
参考资源 - Tokio 官方文档:Cancellation -
tokio_util::sync::CancellationTokenAPI Reference - Google SRE Book: Chapter 21 — Distributed Deadlock

发表评论 取消回复