Rust 异步背压:别让队列替你撒谎
系统处理不过来时,最危险的反应不是报错,而是先把请求悄悄存起来。
队列很容易给人一种安全感:请求来了,先放进去;worker 慢慢处理;入口快速返回。看起来延迟被抹平了,吞吐也稳了。
但如果生产速度持续大于消费速度,队列不是解决问题,只是在延后爆炸。
进入速度 1000/s
处理速度 200/s
每秒堆 800 条
一分钟后堆 48000 条这时候系统表面还活着,实际已经欠了一屁股债。内存涨、延迟涨、重试涨,最后一起崩。
背压要解决的不是“怎么多存一点”,而是“处理不过来时,怎么把压力传回上游”。
这篇接着 超时和取消、优雅关闭 往下写。超时解决“别等太久”,优雅关闭解决“别死太快”,背压解决“别接太多”。
本文代码环境:
# Cargo.toml
[dependencies]
tokio = { version = "1", features = ["full"] }
tower = { version = "0.5", features = ["buffer", "limit", "load-shed", "util"] }队列不是背压,只是缓冲
先看一个很常见的异步写法:
use tokio::sync::mpsc;
#[derive(Debug)]
struct Job(String);
fn spawn_worker() -> mpsc::UnboundedSender<Job> {
let (tx, mut rx) = mpsc::unbounded_channel();
tokio::spawn(async move {
while let Some(job) = rx.recv().await {
process(job).await;
}
});
tx
}这段代码短、好写、能跑。
问题在于 unbounded_channel 没有背压。只要 receiver 没关闭,send 基本都会成功。如果 worker 跟不上,消息会一直堆在内存里。
补一个模拟处理函数:
async fn process(job: Job) {
println!("processing {job:?}");
}这类无界队列适合很窄的场景:消息量有天然上限,或者只是进程内控制信号。拿它承接用户请求、日志事件、数据库写入、远程调用任务,都要非常谨慎。
无界队列默认不该出现在服务主路径上。
有界队列才会把压力传回去
换成有界 channel,语义立刻变了:
use tokio::sync::mpsc;
fn spawn_bounded_worker() -> mpsc::Sender<Job> {
let (tx, mut rx) = mpsc::channel(1024);
tokio::spawn(async move {
while let Some(job) = rx.recv().await {
process(job).await;
}
});
tx
}发送方如果用 send().await:
async fn submit(tx: &mpsc::Sender<Job>, job: Job) -> Result<(), &'static str> {
tx.send(job).await.map_err(|_| "worker stopped")
}队列满了,send().await 会等。
这就是最朴素的背压:下游处理不过来,上游不会继续无限塞。它要么等待,要么超时,要么返回错误。
但这里也有个细节:等待不是永远等。请求路径上最好给它一个时间预算:
use std::time::Duration;
use tokio::time;
async fn submit_with_timeout(tx: &mpsc::Sender<Job>, job: Job) -> Result<(), &'static str> {
time::timeout(Duration::from_millis(50), tx.send(job))
.await
.map_err(|_| "queue full")?
.map_err(|_| "worker stopped")
}这段代码表达的策略是:队列短暂满了可以等 50ms,超过就失败。
背压不是让用户无限等待,而是把系统容量变成明确的等待或拒绝。
try_send 是快速失败,不是偷懒
有些请求不值得排队。
比如遥测事件、在线状态刷新、非关键通知。队列满了,与其让主请求等,不如直接丢掉或降级。
这时用 try_send 更合适:
use tokio::sync::mpsc::error::TrySendError;
fn submit_best_effort(tx: &mpsc::Sender<Job>, job: Job) -> Result<(), &'static str> {
match tx.try_send(job) {
Ok(()) => Ok(()),
Err(TrySendError::Full(_job)) => Err("queue full"),
Err(TrySendError::Closed(_job)) => Err("worker stopped"),
}
}send().await 和 try_send() 没有谁更高级,只有语义不同:
| 写法 | 队列满了怎么办 | 适合场景 |
|---|---|---|
send().await |
等待容量 | 必须处理的任务 |
timeout(send) |
等一小段时间 | 用户请求、RPC 调用 |
try_send |
立刻失败 | 日志、遥测、可丢事件 |
| 无界 channel | 继续堆内存 | 控制信号、小规模内部消息 |
我个人更喜欢先问业务问题:这条消息满了以后,应该等、丢、拒绝,还是转移到持久化队列?
问清楚这个问题,再选 API。
并发不是队列,Semaphore 管的是“正在处理”
队列控制的是“等待处理”的数量。
但很多时候真正要限制的是“同时处理”的数量。比如:
- 同时查数据库的请求不能超过 50
- 同时调用下游支付接口不能超过 20
- 同时跑 CPU 密集任务不能超过核心数
这时候用 Semaphore 更直接:
use std::sync::Arc;
use tokio::sync::Semaphore;
#[derive(Clone)]
struct Limiter {
permits: Arc<Semaphore>,
}
impl Limiter {
fn new(max_in_flight: usize) -> Self {
Self { permits: Arc::new(Semaphore::new(max_in_flight)) }
}
}请求进来先拿 permit:
async fn call_downstream(limiter: Limiter, req: String) -> Result<String, &'static str> {
let _permit = limiter
.permits
.clone()
.acquire_owned()
.await
.map_err(|_| "limiter closed")?;
do_rpc(req).await
}补一个模拟 RPC:
async fn do_rpc(req: String) -> Result<String, &'static str> {
Ok(format!("ok: {req}"))
}_permit 活着的时候,占着一个并发名额。函数返回或出错时,permit drop,名额自动释放。
这里的重点是:Semaphore 限制的是 in-flight,不是 backlog。
如果你在拿 permit 前面再放一个很大的队列,系统照样会堆积。队列和信号量要一起看:一个控制等待区,一个控制执行区。
背压和限流不是一回事
限流和背压很像,但方向不同。
| 机制 | 看什么 | 决策 |
|---|---|---|
| 限流 | 入口速率 | 请求太多,先挡住 |
| 背压 | 下游容量 | 处理不过来,往上游传 |
| 熔断 | 错误状态 | 下游不健康,快速失败 |
| 超时 | 等待时间 | 太慢了,停止等待 |
API 网关里讲过限流:控制进入系统的速率。背压更偏内部链路:某个 worker、channel、连接池、下游服务忙了,调用者应该感知到“现在不能继续塞”。
如果只有限流没有背压,内部某个局部组件还是可能被打爆。
如果只有背压没有限流,入口可能堆满等待中的请求,用户看到的就是越来越慢。
我的偏好是:入口用限流保护整体,内部用背压保护局部。
Tower 里的背压入口是 poll_ready
写 HTTP 服务时,很多背压已经藏在 Tower 抽象里。
Service 不是只有 call,还有一个很关键的 poll_ready:
use std::task::{Context, Poll};
use tower::Service;
fn require_ready<S, Req>(svc: &mut S, cx: &mut Context<'_>) -> Poll<Result<(), S::Error>>
where
S: Service<Req>,
{
svc.poll_ready(cx)
}poll_ready 的语义是:我现在能不能接下一个请求?
如果服务内部队列满了、连接池满了、下游不可用,它可以返回 Pending 或错误。调用方应该先等 ready,再 call。
很多 middleware 写错,就是因为只包了 call,没把 poll_ready 转发给内层服务。这样会绕开内层的背压信号。
正确姿势是这样的:
fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), S::Error>> {
self.inner.poll_ready(cx)
}然后 call 里只做自己的逻辑:
fn call(&mut self, req: Request) -> S::Future {
self.inner.call(req)
}这两段看起来没什么技术含量,但它决定了你的 middleware 有没有尊重下游容量。
在 Tower 体系里,不转发 poll_ready,就等于剪断背压线。
BufferLayer 是缓冲,不是无限容量
Tower 提供了 BufferLayer,可以把请求先放进一个内部队列,再由后台 worker 调 inner service。
它适合削掉一点短时抖动:
use tower::{buffer::Buffer, buffer::BufferLayer, BoxError, Service, ServiceBuilder};
fn build_stack<S, Req>(service: S) -> Buffer<Req, S::Future>
where
S: Service<Req> + Send + 'static,
S::Future: Send + 'static,
S::Error: Into<BoxError> + Send + Sync,
Req: Send + 'static,
{
ServiceBuilder::new()
.layer(BufferLayer::new(128))
.service(service)
}BufferLayer::new(128) 里的 128 很关键。它不是随便写个大数字,而是你愿意承受的排队深度。
buffer 太小,突发流量容易直接失败;buffer 太大,用户请求会在系统里排很久,最后就算成功也没意义。
缓冲只能吸收短抖动,不能解决长期过载。
如果下游持续慢,buffer 迟早满。满了以后,你需要决定:等待、快速失败、降级,还是把任务转成离线处理。
LoadShed 是承认“我现在接不了”
有些系统宁愿快速失败,也不要排队。
比如用户已经在等 HTTP 响应,下游排队 5 秒才处理,成功也没什么价值。与其让请求卡住,不如返回 503,让上游重试或降级。
Tower 里这个策略叫 load shed:
use tower::{load_shed::LoadShed, load_shed::LoadShedLayer, ServiceBuilder};
fn shed_when_busy<S>(service: S) -> LoadShed<S> {
ServiceBuilder::new()
.layer(LoadShedLayer::new())
.service(service)
}load shed 的核心不是“粗暴拒绝”,而是把过载变成显式错误。
排队会让过载看起来没发生,只是延迟越来越高。快速失败会让系统立刻暴露容量不足,调用方也能更早走重试、降级或熔断。
如果请求有明确 deadline,快速失败通常比长时间排队更诚实。
背压指标要直接暴露出来
背压不是写完代码就结束。你要能看见它什么时候触发。
至少应该监控这些:
| 指标 | 看什么 |
|---|---|
| queue_depth | 等待处理的任务数 |
| queue_full_total | 队列满了多少次 |
| in_flight | 正在处理的请求数 |
| permit_wait_ms | 等 Semaphore 花了多久 |
| poll_ready_pending | service 不 ready 的次数 |
| shed_total | 快速失败了多少请求 |
这些指标的价值不只是报警。
它们能帮你判断容量问题在哪一层:入口太猛、队列太小、worker 太慢、数据库连接池太紧,还是下游服务已经不健康。
没有这些指标,背压触发时你只会看到“接口变慢了”。
结论
背压不是某个 API,而是一条容量信号链。
可以直接记这几条:
- 无界队列没有背压,只是把上限交给内存
- 有界队列满了以后,必须明确等、丢、拒绝还是转离线
Semaphore控制 in-flight,不控制 backlogsend().await、try_send、timeout(send)是三种不同业务语义- Tower 里必须尊重
poll_ready - buffer 只能吸收短抖动,不能解决长期过载
- 过载时快速失败往往比长时间排队更可控
系统处理不过来时,要让上游知道,而不是悄悄把压力藏进队列。
队列、Semaphore、poll_ready、load shed 都只是工具。真正的设计问题是:当容量不够时,请求应该在哪里等、等多久、满了以后谁来承担结果。