Pingora 限流实战:基于 pingora-limits 的 Rate 速率限制器快速上手指南
2026/9/10 23:22:51 网站建设 项目流程

Pingora 限流实战:基于 pingora-limits 的 Rate 速率限制器快速上手指南

【免费下载链接】pingoraA library for building fast, reliable and evolvable network services.项目地址: https://gitcode.com/GitHub_Trending/pi/pingora

导读

本文围绕 Pingora 框架官方用户指南 rate_limiter.md 展开,完整讲解如何利用pingora-limitscrate 提供的Rate类型,在ProxyHttprequest_filter阶段按请求头(如appid)为每个客户端实现每秒请求数(QPS)限流,并在超限时返回429 Too Many Requests及标准X-Rate-Limit-*响应头。读完本文,你将掌握Rate的核心 API(observe/rate/rate_with)、其底层"双槽计数器 + Count-Min Sketch"实现原理,以及一个可直接cargo run运行的完整反向代理限流示例。


一、限流方案概览

Pingora 提供了独立的限流 cratepingora-limits(源码位于 pingora-limits/src/lib.rs),它包含三个模块:

  • estimator:无锁 Count-Min Sketch 频率估计器,是其余模块的底层存储;
  • rateRate类型,用于估计一段时间窗口内事件的发生频率;
  • inflightInflight类型,用于估计某一时刻正在发生的事件数(并发量)。

官方指南给出的应用场景是:以请求头appid区分调用方,为每个appid单独维护一个速率限制器,限制其每秒请求数不超过阈值。整体流程只有三步:

  1. Cargo.toml中加入依赖;
  2. 声明一个全局的Rate限流器(按 key 如appid区分不同客户端);
  3. 覆写ProxyHttptrait 的request_filter方法,实现计数、判断与429响应。

这一设计把"限流"完全嵌入代理请求处理管线,无需额外的外部服务(如 Redis),也无需加锁,非常适合在网关/反向代理中做轻量级、分布式的单机限流。

二、添加依赖

在应用的Cargo.toml中加入以下依赖(与官方指南一致):

async-trait="0.1" pingora = { version = "0.8", features = [ "lb", "openssl" ] } pingora-limits = "0.8.0" once_cell = "1.19.0"

说明:

  • pingoralbfeature 提供LoadBalancerRoundRobinTcpHealthCheck等负载均衡组件;opensslfeature 提供 TLS 支持。当前仓库中pingora-limits的版本正是0.8.0(见 pingora-limits/Cargo.toml)。
  • pingora-limits本身依赖极少,生产依赖只有ahash,其正式描述为 "A library for rate limiting and event frequency estimation"。
  • once_cell用于以Lazy方式声明全局静态限流器。

三、核心实现:基于Rate的按客户端限流

下面解析官方指南示例 rate_limiter.rs(该文件可直接运行,仓库路径见后文"完整示例"一节)。

3.1 全局限流器与阈值

use once_cell::sync::Lazy; use pingora_limits::rate::Rate; use std::time::Duration; // Rate limiter static RATE_LIMITER: Lazy<Rate> = Lazy::new(|| Rate::new(Duration::from_secs(1))); // max request per second per client static MAX_REQ_PER_SEC: isize = 1;

Rate::new(Duration)创建以指定时长为窗口的速率限制器,RATE_LIMITER是进程级全局单例,所有请求共享。MAX_REQ_PER_SEC = 1表示每个appid每秒最多允许 1 个请求,你可以按业务需要调整该阈值。

3.2 从请求头提取客户端标识

pub struct LB(Arc<LoadBalancer<RoundRobin>>); impl LB { pub fn get_request_appid(&self, session: &mut Session) -> Option<String> { match session .req_header() .headers .get("appid") .map(|v| v.to_str()) { None => None, Some(v) => match v { Ok(v) => Some(v.to_string()), Err(_) => None, }, } } }

这里直接从Session的请求头中读取appid字符串;如果客户端未携带appid,则返回None,后续将跳过限流。appid只是示例 key,实际生产中可以替换为 IP、用户 ID、API Key 或任意可Hash的类型。

3.3 在request_filter中完成计数与拦截

#[async_trait] impl ProxyHttp for LB { type CTX = (); fn new_ctx(&self) {} // ... upstream_peer / upstream_request_filter 略 ... async fn request_filter(&self, session: &mut Session, _ctx: &mut Self::CTX) -> Result<bool> where Self::CTX: Send + Sync, { let appid = match self.get_request_appid(session) { None => return Ok(false), // no client appid found, skip rate limiting Some(addr) => addr, }; // retrieve the current window requests let curr_window_requests = RATE_LIMITER.observe(&appid, 1); if curr_window_requests > MAX_REQ_PER_SEC { // rate limited, return 429 let mut header = ResponseHeader::build(429, None).unwrap(); header .insert_header("X-Rate-Limit-Limit", MAX_REQ_PER_SEC.to_string()) .unwrap(); header.insert_header("X-Rate-Limit-Remaining", "0").unwrap(); header.insert_header("X-Rate-Limit-Reset", "1").unwrap(); session.set_keepalive(None); session .write_response_header(Box::new(header), true) .await?; return Ok(true); } Ok(false) } }

关键点:

  • RATE_LIMITER.observe(&appid, 1)为当前appid在当前窗口内增加 1 个事件,并返回该窗口内累计的(估计)事件数。返回值与阈值比较即可判断是否超限。
  • 超限时构造429响应,附带三个标准的限流响应头:
    • X-Rate-Limit-Limit: 1—— 窗口内允许的最大请求数;
    • X-Rate-Limit-Remaining: 0—— 当前剩余配额;
    • X-Rate-Limit-Reset: 1—— 距窗口重置的秒数(本例窗口为 1s)。
  • session.set_keepalive(None)关闭 keepalive,write_response_header(.., true)直接写回响应并结束本次会话。
  • 返回值语义:Ok(true)表示代理已自行处理完该请求(终止后续处理),Ok(false)表示继续走正常的代理流程(转发到上游)。

关于request_filter的位置:它是ProxyHttptrait 定义的请求处理钩子,官方文档注释明确指出该阶段用于"解析、校验、限流、访问控制或直接返回响应"(见 proxy_trait.rs)。默认实现为空并返回Ok(false)。如果你希望限流逻辑在其他模块之前执行,还可以考虑early_request_filter,但按注释建议,能放在request_filter就放在这里,以便同样受其他模块的访问控制保护。

3.4 主函数装配

fn main() { let mut server = Server::new(Some(Opt::default())).unwrap(); server.bootstrap(); let mut upstreams = LoadBalancer::try_from_iter(["1.1.1.1:443", "1.0.0.1:443"]).unwrap(); // Set health check let hc = TcpHealthCheck::new(); upstreams.set_health_check(hc); upstreams.health_check_frequency = Some(Duration::from_secs(1)); // Set background service let background = background_service("health check", upstreams); let upstreams = background.task(); // Set load balancer let mut lb = http_proxy_service(&server.configuration, LB(upstreams)); lb.add_tcp("0.0.0.0:6188"); server.add_service(background); server.add_service(lb); server.run_forever(); }

主函数监听0.0.0.0:6188,将流量轮询转发到1.1.1.1:443/1.0.0.1:443(Cloudflare 公共 DNS 的 DoH 端点),并设置 SNI 为one.one.one.one(见upstream_peerupstream_request_filter)。健康检查以 1 秒为周期后台运行。这些基础设施代码表明限流器可以无缝嵌入一个完整的、带负载均衡与健康检查的生产级代理服务。

四、Rate的底层实现原理

Rate定义于 pingora-limits/src/rate.rs。它采用双槽(double-buffer)计数器 + 无锁 Count-Min Sketch设计,能够在多线程高并发下以极低成本完成计数。

4.1 双槽窗口切换

pub struct Rate { red_slot: Estimator, blue_slot: Estimator, red_or_blue: AtomicBool, // true: the current slot is red, otherwise blue start: Instant, reset_interval_ms: u64, last_reset_time: AtomicU64, interval: Duration, }
  • red_slotblue_slot两个Estimator轮流充当"当前窗口"与"上一个窗口":当前窗口用于收集事件,上一个已完成的窗口用于报告速率。
  • 默认构造参数为HASHES = 4SLOTS = 1024(rate.rs),可通过Rate::new_with_estimator_config(interval, hashes, slots)自定义;注释提醒:如果窗口较短、key 基数较低,SLOTS可以调小。
  • 窗口切换逻辑在maybe_reset()中:用last_reset_time记录上次重置时刻,通过compare_exchange原子地完成"清空旧槽 → 翻转red_or_blue标志",从而避免多线程竞争时重复重置;若距上次重置已超过两个窗口,则两个槽都会被清空(rate.rs)。

4.2 核心 API

方法作用返回
Rate::new(interval)创建窗口为interval的限流器Rate
observe(&key, events)key在当前窗口累加events个事件当前窗口内累计(估计)事件数
rate(&key)返回key最近一个完整窗口的平均速率(每秒事件数)f64
rate_with(&key, calc_fn)用自定义闭包计算速率闭包返回类型
new_with_estimator_config(interval, hashes, slots)自定义哈希数与槽位数Rate

其中rate()的公式为:上一个窗口计数 * 1000 / 窗口毫秒数,即"每秒事件数"(rate.rs)。当距上次重置超过两个窗口(无新事件)时,直接返回 0 作为短路优化。

4.3 平滑速率估计与自定义计算

RateComponents结构体向自定义速率函数暴露四项信息(rate.rs):

  • prev_samples:上一完整窗口的样本数;
  • curr_samples:当前进行中窗口的样本数;
  • interval:窗口时长;
  • current_interval_fraction:采样点处于当前窗口的进度比例(0..1),例如窗口 10s、在第 2 秒采样时为 0.2。

内置的PROPORTIONAL_RATE_ESTIMATE_CALC_FN使用线性插值在上一窗口与当前窗口之间加权,得到"过去interval时间内"的平滑速率估计:

let weighted_count = prev * (1. - interval_fraction) + curr; weighted_count / interval_secs

该思路源自 Cloudflare 的博客文章《Counting things a lot of different things》——用少量内存近似统计海量 key 的频率。如果你想要"90% 当前 + 10% 历史"之类的自定义策略,可以通过rate_with传入自己的闭包(参考 rate.rs 中的test_observe_rate_custom_90_10测试)。

4.4 底层存储:无锁 Count-Min Sketch

Estimator(estimator.rs)是标准的 Count-Min Sketch:hashes行、每行slotsAtomicIsize计数器,每个 key 经 4 路独立哈希(ahash::RandomState)落到每行的一个槽位,读取时取所有行计数的最小值作为频率估计。插入与查询均为无锁原子操作(fetch_add/loadOrdering::Relaxed),时间复杂度 O(h)、空间复杂度 O(h×n)。需要注意:计数溢出时可能返回负数,调用方需自行处理(estimator.rs)。

因此Rateobserve返回的是估计值而非精确值——这是以极低内存代价换取高吞吐的典型取舍。若需要精确并发计数,pingora-limits还提供了Inflight类型(inflight.rs),它基于同样的Estimator,通过 RAIIGuard在离开作用域时自动decr计数。

五、测试与验证

5.1 用 curl 验证限流效果

官方指南给出了完整验证步骤:运行程序后,用以下命令连续发送携带appid的请求:

curl localhost:6188 -H "appid:1" -v

由于MAX_REQ_PER_SEC = 1第一个请求应成功转发;随后 1 秒窗口内的后续请求应收到429

* Trying 127.0.0.1:6188... * Connected to localhost (127.0.0.1) port 6188 (#0) > GET / HTTP/1.1 > Host: localhost:6188 > User-Agent: curl/7.88.1 > Accept: */* > appid:1 > < HTTP/1.1 429 Too Many Requests < X-Rate-Limit-Limit: 1 < X-Rate-Limit-Remaining: 0 < X-Rate-Limit-Reset: 1 < Date: Sun, 14 Jul 2024 20:29:02 GMT < Connection: close < * Closing connection 0

可以尝试:

  • 更换appid值(如-H "appid:2"),验证限流是按客户端独立统计的;
  • 不携带appid,验证请求会跳过限流正常转发;
  • 等待 1 秒后再发请求,验证窗口重置后重新放行。

5.2 单元测试中的速率语义

Rate自带单元测试(rate.rs),直观展示了窗口语义:

  • 窗口内observe累加返回累计值;未跨窗口时rate()为 0;
  • 跨过 1 个窗口后,rate()返回上一窗口的事件数(即每秒速率);
  • 超过 2 个窗口没有事件,rate()归零。

由于测试中大量使用真实sleep,测试代码以宽容误差(epsilon = 0.15)做近似断言。这些测试同时覆盖了PROPORTIONAL_RATE_ESTIMATE_CALC_FN的插值行为,可作为理解窗口切换逻辑的补充材料。

六、完整示例运行

仓库已在 pingora-proxy/examples/rate_limiter.rs 内置了与本文完全一致的可运行示例(即官方指南正文代码的完整版,包含use导入与 load balancer 相关类型)。进入pingora-proxycrate 目录后执行:

cargo run --example rate_limiter

程序启动后将监听0.0.0.0:6188,然后按上一节的 curl 命令即可验证限流效果。

七、适用前提与扩展方向

  • 适用前提Rate是单进程内的内存限流器。对于多副本部署或跨机器限流,需要在各实例间同步计数(如借助共享存储或消息队列),本文方案不直接适用;且Rate基于 Count-Min Sketch 做频率估计,计数为近似值,追求精确计数的场景请谨慎评估。
  • 阈值与窗口:通过调整Rate::newDurationMAX_REQ_PER_SEC,可实现"每秒 X 次""每分钟 Y 次"等不同粒度;需要更精细的内存/精度权衡时,可使用Rate::new_with_estimator_config调整哈希数与槽位数。
  • 限流依据:示例使用appid请求头,实际可改为客户端 IP、Cookie、认证后的用户 ID 等任意可哈希键;未携带标识的请求可配置为放行或默认配额。
  • 响应增强X-Rate-Limit-*响应头遵循常见的限流响应约定,可在此基础上补充Retry-After头或 JSON 错误体,便于客户端做退避重试。

通过本文,你已掌握 Pingora 框架内嵌限流能力从依赖、编码、原理到验证的完整闭环,可以直接将其移植到自己的网关或代理服务中。

【免费下载链接】pingoraA library for building fast, reliable and evolvable network services.项目地址: https://gitcode.com/GitHub_Trending/pi/pingora

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询