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::CancellationToken API Reference - Google SRE Book: Chapter 21 — Distributed Deadlock

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部