Tokio 在生产环境:日处理千万请求的异步服务怎么配置才稳妥 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: HashMapString, ArcSemaphore, } 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, ) - Resulttokio::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): StateArcRateLimiter, ) - Resultimpl 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本身它持有ArcSemaphore这意味着你可以在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: ArcNotify, } impl GracefulShutdown { pub fn new() - Self { Self { notify_shutdown: Arc::new(Notify::new()), } } /// 监听系统信号并触发关闭 pub async fn listen(self) { let ctrl_c async { // 捕获 CtrlC (SIGINT) signal::ctrl_c() .await .expect(无法安装 CtrlC handler); }; let terminate async { // 捕获 SIGTERMK8s pod 删除时发送的信号 signal::unix::signal(signal::unix::SignalKind::terminate()) .expect(无法安装 SIGTERM handler) .recv() .await; }; // 等待任意一个信号到达 tokio::select! { _ ctrl_c { println!(收到 CtrlC开始优雅关闭...); } _ terminate { println!(收到 SIGTERM开始优雅关闭...); } } // 通知所有等待的任务开始关闭 self.notify_shutdown.notify_waiters(); } /// 返回一个用于等待关闭通知的 Future pub fn shutdown_signal(self) - impl std::future::FutureOutput () { 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 留余量给 sidecarmax_blocking_threads 根据 spawn_blocking 用量调整Semaphore 做 per-endpoint 限流防止单点过载拖垮全局通过tokio::signal GracefulShutdown 实现安全的重启流程。日处理千万请求不是一个很难的数字但要做到稳妥确实需要把这些细节都处理好。从 Go 迁到 Rust 最大的感受不是 QPS 升了多少而是内存稳了多少——迁移后内存从平均 800MB 降到了 120MB而且再也没碰到过 OOM。省下来的内存都给业务用了。

本月热点