Skip to main content

Rust 高并发实战(3):Axum、Semaphore 与生产级有界服务

Rainy
雨落无声,代码成诗 —— 致力于技术与艺术的极致平衡
Rainy
7 MIN READ... VIEWS

async 让等待更便宜,Semaphore 决定允许多少等待同时存在。

本期实现 /aggregate:每个请求并发访问多个下游,但入口和依赖都有独立硬上限。

架构图:Axum 请求经过哪些并发边界

Tower 入口限制保护整个 route;三个依赖 semaphore 分别保护真实故障域;连接池是最后一道物理边界。若 semaphore 大于连接池,task 会在两处排队,必须分别记录等待时间。

一、依赖与状态

[dependencies]
axum = "0.8"
futures-util = "0.3"
serde = { version = "1", features = ["derive"] }
tokio = { version = "1", features = ["macros", "rt-multi-thread", "signal", "sync", "time"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] }
#[derive(Clone)]
struct AppState {
requests: Arc<Semaphore>,
database: Arc<Semaphore>,
rpc: Arc<Semaphore>,
accepted: Arc<AtomicU64>,
rejected: Arc<AtomicU64>,
}

入口池保护进程;DB/RPC 池按故障域隔离。许可数必须与实际连接池、下游副本和全服务副本数一起计算。

二、准入与总 deadline

#[tracing::instrument(skip(state))]
async fn aggregate(
State(state): State<AppState>,
) -> Result<Json<AggregateResponse>, ApiError> {
let request_permit = tokio::time::timeout(
Duration::from_millis(20),
state.requests.clone().acquire_owned(),
)
.await
.map_err(|_| {
state.rejected.fetch_add(1, Ordering::Relaxed);
ApiError::Overloaded
})?
.map_err(|_| ApiError::ShuttingDown)?;

state.accepted.fetch_add(1, Ordering::Relaxed);
let work = async {
let database = query_database(state.clone());
let rpc = call_rpc(state.clone());
let cache = read_cache(state.clone());
let (database, rpc, cache) = tokio::try_join!(database, rpc, cache)?;
Ok(AggregateResponse { database, rpc, cache })
};

let response = tokio::time::timeout(Duration::from_millis(250), work)
.await
.map_err(|_| ApiError::Deadline)??;
drop(request_permit);
Ok(Json(response))
}

OwnedSemaphorePermit 通过 RAII 在所有返回路径释放。不要提前 drop,也不要 spawn 一个仅为“持有许可”的 detached task。

三、下游许可与剩余预算

async fn query_database(state: AppState) -> Result<DatabaseValue, ApiError> {
let _permit = tokio::time::timeout(
Duration::from_millis(15),
state.database.acquire_owned(),
)
.await
.map_err(|_| ApiError::Overloaded)?
.map_err(|_| ApiError::ShuttingDown)?;

tokio::time::timeout(Duration::from_millis(80), execute_query())
.await
.map_err(|_| ApiError::DependencyTimeout)?
}

生产代码最好传递一个绝对 deadline 或剩余预算,而不是每层重新获得完整 250 ms。分别记录 semaphore queue wait 与真正 I/O duration。

四、错误分类

类型HTTP是否重试监控含义
入口过载429按 Retry-After + jitter容量饱和
停机503可换实例部署/摘流
总 deadline504取决于幂等与预算长尾
下游超时502/504受重试预算限制依赖故障
参数错误400调用方问题

不要全映射成 500,否则客户端策略、SLO 和容量指标都会失真。

五、优雅停机

let app = Router::new()
.route("/aggregate", get(aggregate))
.with_state(state);
let listener = TcpListener::bind("0.0.0.0:8080").await?;

axum::serve(listener, app)
.with_graceful_shutdown(shutdown_signal())
.await?;

部署顺序:readiness 失败 → 等摘流传播 → 触发 graceful shutdown → 限时 drain。后台 task 需要单独的 CancellationToken/JoinSet owner,不能假设 Axum 会收割所有 detached task。

六、Tower Layer 的顺序

Axum 基于 Tower。timeout、concurrency limit、load shed、trace 等 layer 的嵌套顺序会改变语义:timeout 是否包含排队?trace 是否记录被拒绝请求?load shed 在等待前还是后触发?

let app = Router::new()
.route("/aggregate", get(aggregate))
.layer(TraceLayer::new_for_http())
.layer(TimeoutLayer::new(Duration::from_millis(300)))
.with_state(state);

不要同时在 Tower 和 handler 放两套意义不明的同级并发限制。入口 layer 保护整个 route,依赖 semaphore 保护具体下游;分别命名指标。

七、请求体、响应与内存上限

高并发下每个 10 MiB body × 500 in-flight 就不可接受。对 body 设置硬限制,流式读取时保留总大小和 deadline。聚合响应也不能无界收集:限制 fan-out、单结果大小和总结果大小。

序列化是 CPU 与分配工作,超大 JSON 会占用 runtime worker。对大响应考虑分页/流式协议,CPU 密集编码按 profile 决定是否隔离。

八、客户端与连接池

复用 HTTP/DB client;每请求创建 client 会破坏连接复用。连接池容量小于 semaphore 会形成二次排队,容量大于下游预算又会把过载推给依赖。

每个阶段记录:permit wait、pool acquisition、DNS/connect/TLS、TTFB、body read。超时错误保留阶段信息,不能全部变成 elapsed

九、完整关闭协议

readiness=false
↓ 等待 LB 传播
停止 listener / 拒绝新任务

CancellationToken 通知后台任务

JoinSet/TaskTracker 等待 drain
↓ 超过 shutdown deadline
记录未完成工作并强退

spawn_blocking 已开始的 closure 不会因 abort 停止,必须在停机预算中单独考虑。消息消费者要停止拉取新消息,再完成或归还当前消息。

十、指标与基数控制

Prometheus label 不放 request ID、用户 ID、原始 URL 或错误全文。route 使用模板 /users/:id;错误按稳定 class 分类,详细上下文放 tracing/log。

推荐直方图:admission_wait_secondsdependency_wait_secondsdependency_duration_secondsrequest_duration_seconds。计数:accepted/rejected/timeout/cancelled/result_class。

十一、慢客户端与连接边界

应用 deadline 不一定覆盖客户端缓慢上传 header/body。生产入口通常还需代理/Hyper 层的 header、body、keep-alive 和连接限制。代理 timeout 与应用 timeout 要分层命名,避免代理 504 时应用仍继续工作。

对 SSE/WebSocket/流式响应不能使用普通短 Write timeout,要设置心跳、idle timeout、每连接 buffer 与全局连接上限,并在 shutdown 时通知长连接关闭。

十二、部分结果策略

聚合服务要把依赖分为 required 和 optional。required 失败则整体失败;optional 失败可返回 stale/缺省值,同时在 response metadata 标记 degraded:

#[derive(Serialize)]
struct AggregateResponse {
user: User,
recommendations: Option<Vec<Item>>,
degraded: Vec<&'static str>,
}

降级成功不能完全从错误率消失,应有独立 degraded 指标和 SLO,否则系统会长期带病运行。

12.1 把剩余 deadline 变成显式值

#[derive(Clone, Copy)]
struct Deadline(Instant);

impl Deadline {
fn remaining(self) -> Result<Duration, ApiError> {
self.0.checked_duration_since(Instant::now())
.filter(|left| !left.is_zero())
.ok_or(ApiError::Deadline)
}

async fn run<T>(self, future: impl Future<Output = T>) -> Result<T, ApiError> {
tokio::time::timeout(self.remaining()?, future)
.await
.map_err(|_| ApiError::Deadline)
}
}

入口创建绝对 deadline,子层只读取剩余时间。数据库可能最多拿 80 ms,但如果入口只剩 35 ms,实际 timeout 是 35 ms。optional 调用还要保留编码/写回余量,不应耗尽最后一毫秒。

12.2 双重排队实验

把 DB semaphore 设置 100、连接池设置 10,在 500 并发下观察:90 个 task 已通过应用保护却堵在连接池。然后把 semaphore 对齐到 10~12,比较 pool wait、task 数、RSS 和 P99。目标不是证明 10 最快,而是让排队发生在可观测、可超时的单一位置。

十三、本期验收

  • 分别压满入口、DB、RPC semaphore,确认故障隔离;
  • 注入 1 s 下游延迟,确认 250 ms 后 handler 和子 Future 被取消;
  • 检查许可、连接和 task 数能回落;
  • 滚动停机时无明显 5xx 峰值;
  • 记录 accepted、rejected、queue wait、dependency duration 和 result class。

上一篇:Rust-2:Future 与 Tokio · 下一篇:Rust-4:锁、Channel、Atomic 与取消安全

参考资料

Logo
RainLib

Exploring the frontiers of technology, design, and distributed systems. Building tools for the future developers.

Suggestions & Feedback

© 2026 RainLib. Built for the Future.
All rights reserved.
System Normal