返回文章列表

Rust 异步背压:别让队列替你撒谎

472·4 分钟阅读
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().awaittry_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,不控制 backlog
  • send().awaittry_sendtimeout(send) 是三种不同业务语义
  • Tower 里必须尊重 poll_ready
  • buffer 只能吸收短抖动,不能解决长期过载
  • 过载时快速失败往往比长时间排队更可控

系统处理不过来时,要让上游知道,而不是悄悄把压力藏进队列。

队列、Semaphore、poll_ready、load shed 都只是工具。真正的设计问题是:当容量不够时,请求应该在哪里等、等多久、满了以后谁来承担结果。