Tokio 在生产环境:日处理千万请求的异步服务怎么配置才稳妥
一、从"能用"到"稳用",Tokio 默认配置的陷阱
刚开始学 Rust 异步编程的时候,一个#[tokio::main]就能启动一个异步运行时,一切都显得那么自然。我也天真地以为生产环境也能这么搞——直到压测的时候发现,QPS 到 5000 以后延迟开始飙升,到 8000 就直接 OOM 了。
问题在哪?Tokio 的默认配置是为开发环境设计的。默认的 worker 线程数等于 CPU 核心数,这在 4 核的笔记本上没问题,但在 32 核的生产服务器上反而会成为瓶颈——因为太多线程竞争同一个任务队列会产生大量上下文切换开销。而且默认的max_blocking_threads只有 512,高并发时可能会爆。
图上把最常见的三个陷阱和对应的解决方案都画出来了。核心原则就是:不要把 Tokio 当作黑盒——你越了解它怎么调度任务,就越知道该在哪里加保护。
二、Runtime 调优:这才是生产环境的第一课
下面是我在实际迁移中总结的 runtime 配置。注意worker_threads我设成了num_cpus - 2,因为服务器上还跑着 Envoy sidecar 和监控 agent,需要给它们留出 CPU 时间。max_blocking_threads调到 2048 是因为我们有一些 DB 查询走的是spawn_blocking,高峰时需要容纳大量等待数据库返回的线程。
use tokio::runtime::{Builder, Runtime}; use std::time::Duration; /// 构建生产环境的 Tokio runtime /// 区别于默认的 #[tokio:main],手动配置所有关键参数 fn build_production_runtime() -> Runtime { // 获取 CPU 核心数 let num_cpus = num_cpus::get(); // 预留 2 个核心给 OS 和 sidecar 进程 let worker_threads = (num_cpus.saturating_sub(2)).max(2); Builder::new_multi_thread() // 设置 worker 线程数,不超出 CPU 核心数 .worker_threads(worker_threads) // 每个 worker 线程可同时处理的最大任务数 // 设为 256 适合 IO 密集型,CPU 密集型建议 32 .max_blocking_threads(2048) // 启用 IO 驱动,epoll/kqueue 的事件循环 .enable_io() // 启用时间驱动,tokio::time::sleep 等功能需要 .enable_time() // 全局任务队列间隔(影响公平性) // 值越小越公平,但吞吐量可能略微下降 .global_queue_interval(61) // 默认值 61,一般不需要改 // 线程名称前缀(方便用 htop / perf 定位问题线程) .thread_name("tokio-gateway-worker") // 线程栈大小(默认 2MB,这里不变) // 如果你的服务递归深,可以适当调大 .build() .expect("构建 Tokio runtime 失败") }global_queue_interval这个参数我第一次见到时完全不懂它的含义,查了源码才明白:Tokio 的多线程调度器有一个全局任务队列和每个 worker 的本地队列。global_queue_interval控制 worker 每隔多少轮去全局队列偷一次任务。值设小一点可以让任务分布更均匀,但会增加全局队列的锁竞争。61 是默认值,对大多数场景是合理的,我个人习惯不动它。
三、背压与限流:服务不崩的底线
runtime 配好了,接下来是最容易被忽略但最致命的问题——背压(backpressure)。异步服务的特点是可以同时接受大量连接,但如果下游处理不过来,连接和内存就会无限堆积直到 OOM。
我用tokio::sync::Semaphore实现了一个简单的 per-endpoint 限流器。每个 API 端点独立配置最大并发数,这样即使某个 endpoint 被打爆了,也不会影响其他正常的接口。
use tokio::sync::Semaphore; use std::sync::Arc; use std::collections::HashMap; /// 端点级别限流器 /// 每个 API 端点有独立的信号量,互不影响 pub struct RateLimiter { /// 端点名 -> 信号量(最大并发许可数) limits: HashMap<String, Arc<Semaphore>>, } impl RateLimiter { pub fn new() -> Self { let mut limits = HashMap::new(); // 查询类接口:并发量大,设 500 limits.insert("/api/query".to_string(), Arc::new(Semaphore::new(500))); // 写入类接口:并发量小,设 100(保护 DB 连接池) limits.insert("/api/write".to_string(), Arc::new(Semaphore::new(100))); // 管理类接口:几乎无并发,设 10 limits.insert("/api/admin".to_string(), Arc::new(Semaphore::new(10))); Self { limits } } /// 获取指定端点的许可(带超时) /// 如果 5 秒内获取不到许可,说明该端点已经过载 pub async fn acquire( &self, endpoint: &str, ) -> Result<tokio::sync::OwnedSemaphorePermit, String> { // 获取该端点对应的信号量 let sem = self.limits .get(endpoint) .ok_or_else(|| format!("未知端点: {}", endpoint))? .clone(); // 尝试在 5 秒内获取许可 match tokio::time::timeout( Duration::from_secs(5), sem.acquire_owned(), // acquire_owned 返回 OwnedSemaphorePermit,生命周期独立 ).await { Ok(Ok(permit)) => Ok(permit), Ok(Err(_)) => Err("信号量已关闭".to_string()), // 正常情况下不会发生 Err(_) => Err(format!("端点 {} 过载,获取许可超时", endpoint)), } } } // ===== 在 axum handler 中使用限流器 ===== use axum::{extract::State, http::StatusCode, response::IntoResponse}; /// 模拟的查询接口 handler async fn query_handler( State(limiter): State<Arc<RateLimiter>>, ) -> Result<impl IntoResponse, StatusCode> { // 尝试获取 /api/query 端点的许可 // 如果获取失败,返回 503 Service Unavailable let _permit = limiter.acquire("/api/query") .await .map_err(|_| StatusCode::SERVICE_UNAVAILABLE)?; // _permit 在函数结束时自动释放回信号量 // 这是 Rust RAII 的经典模式 // 模拟业务处理 let result = process_query().await; Ok(axum::Json(result)) }这里用OwnedSemaphorePermit而不是SemaphorePermit是一个细节。OwnedSemaphorePermit的生命周期独立于Semaphore本身(它持有Arc<Semaphore>),这意味着你可以在tokio::spawn中将 permit 传递到其他任务中去,非常灵活。我刚开始没注意到这个区别,结果在 spawn 时被生命周期错误卡了好久。
四、优雅关闭与连接排空
服务配得再好,总归要重启。如果重启时直接 kill 进程,正在处理的请求就会中断,用户看到的就是 502。Tokio 提供了tokio::signal来捕获系统信号,配合GracefulShutdown可以实现优雅关闭。
use tokio::signal; use tokio::sync::Notify; use std::sync::Arc; use axum::Router; /// 优雅关闭管理器 pub struct GracefulShutdown { /// 通知所有任务开始关闭 notify_shutdown: Arc<Notify>, } impl GracefulShutdown { pub fn new() -> Self { Self { notify_shutdown: Arc::new(Notify::new()), } } /// 监听系统信号并触发关闭 pub async fn listen(&self) { let ctrl_c = async { // 捕获 Ctrl+C (SIGINT) signal::ctrl_c() .await .expect("无法安装 Ctrl+C handler"); }; let terminate = async { // 捕获 SIGTERM(K8s pod 删除时发送的信号) signal::unix::signal(signal::unix::SignalKind::terminate()) .expect("无法安装 SIGTERM handler") .recv() .await; }; // 等待任意一个信号到达 tokio::select! { _ = ctrl_c => { println!("收到 Ctrl+C,开始优雅关闭..."); } _ = terminate => { println!("收到 SIGTERM,开始优雅关闭..."); } } // 通知所有等待的任务开始关闭 self.notify_shutdown.notify_waiters(); } /// 返回一个用于等待关闭通知的 Future pub fn shutdown_signal(&self) -> impl std::future::Future<Output = ()> { let notify = self.notify_shutdown.clone(); async move { notify.notified().await; } } } /// 启动 HTTP 服务(带优雅关闭) async fn start_server_with_graceful_shutdown() { let shutdown = GracefulShutdown::new(); let app = Router::new() .route("/api/query", axum::routing::get(query_handler)) .with_state(/* ... */); let listener = tokio::net::TcpListener::bind("0.0.0.0:8080") .await .unwrap(); println!("服务启动在 0.0.0.0:8080"); // 启动监听信号的 task let shutdown_listener = tokio::spawn(async move { shutdown.listen().await; }); // axum::serve 自带 graceful shutdown 支持 axum::serve(listener, app) .with_graceful_shutdown(async { // 等待关闭信号 shutdown_listener.await.ok(); println!("开始排空连接,等待所有请求完成..."); }) .await .unwrap(); }优雅关闭的关键两步:停止接收新连接 + 等待现有请求完成。axum::serve().with_graceful_shutdown()已经帮我们做了第一步,但第二步需要你在 handler 中检查关闭信号。如果你的请求处理时间很长(比如批量导出),建议在 handler 中也监听关闭通知,及时中断超长任务。
五、总结
这篇文章复盘了 Tokio 从开发环境到生产环境的配置要点:runtime 参数需要根据服务器实际情况调优(worker_threads 留余量给 sidecar,max_blocking_threads 根据 spawn_blocking 用量调整),Semaphore 做 per-endpoint 限流防止单点过载拖垮全局,通过tokio::signal+ GracefulShutdown 实现安全的重启流程。
日处理千万请求不是一个很难的数字,但要做到"稳妥"确实需要把这些细节都处理好。从 Go 迁到 Rust 最大的感受不是 QPS 升了多少,而是内存稳了多少——迁移后内存从平均 800MB 降到了 120MB,而且再也没碰到过 OOM。省下来的内存都给业务用了。