十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

Pingora 限流实战:基于 pingora-limits 的 Rate 速率限制器快速上手指南 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类型在ProxyHttp的request_filter阶段按请求头如appid为每个客户端实现每秒请求数QPS限流并在超限时返回429 Too Many Requests及标准X-Rate-Limit-*响应头。读完本文你将掌握Rate的核心 APIobserve/rate/rate_with、其底层双槽计数器 Count-Min Sketch实现原理以及一个可直接cargo run运行的完整反向代理限流示例。一、限流方案概览Pingora 提供了独立的限流 cratepingora-limits源码位于 pingora-limits/src/lib.rs它包含三个模块estimator无锁 Count-Min Sketch 频率估计器是其余模块的底层存储rateRate类型用于估计一段时间窗口内事件的发生频率inflightInflight类型用于估计某一时刻正在发生的事件数并发量。官方指南给出的应用场景是以请求头appid区分调用方为每个appid单独维护一个速率限制器限制其每秒请求数不超过阈值。整体流程只有三步在Cargo.toml中加入依赖声明一个全局的Rate限流器按 key 如appid区分不同客户端覆写ProxyHttptrait 的request_filter方法实现计数、判断与429响应。这一设计把限流完全嵌入代理请求处理管线无需额外的外部服务如 Redis也无需加锁非常适合在网关/反向代理中做轻量级、分布式的单机限流。二、添加依赖在应用的Cargo.toml中加入以下依赖与官方指南一致async-trait0.1 pingora { version 0.8, features [ lb, openssl ] } pingora-limits 0.8.0 once_cell 1.19.0说明pingora的lbfeature 提供LoadBalancer、RoundRobin、TcpHealthCheck等负载均衡组件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: LazyRate 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(ArcLoadBalancerRoundRobin); impl LB { pub fn get_request_appid(self, session: mut Session) - OptionString { 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) - Resultbool 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)关闭 keepalivewrite_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:443Cloudflare 公共 DNS 的 DoH 端点并设置 SNI 为one.one.one.one见upstream_peer与upstream_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_slot与blue_slot两个Estimator轮流充当当前窗口与上一个窗口当前窗口用于收集事件上一个已完成的窗口用于报告速率。默认构造参数为HASHES 4、SLOTS 1024rate.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的限流器Rateobserve(key, events)为key在当前窗口累加events个事件当前窗口内累计估计事件数rate(key)返回key最近一个完整窗口的平均速率每秒事件数f64rate_with(key, calc_fn)用自定义闭包计算速率闭包返回类型new_with_estimator_config(interval, hashes, slots)自定义哈希数与槽位数Rate其中rate()的公式为上一个窗口计数 * 1000 / 窗口毫秒数即每秒事件数rate.rs。当距上次重置超过两个窗口无新事件时直接返回 0 作为短路优化。4.3 平滑速率估计与自定义计算RateComponents结构体向自定义速率函数暴露四项信息rate.rsprev_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 SketchEstimatorestimator.rs是标准的 Count-Min Sketchhashes行、每行slots个AtomicIsize计数器每个 key 经 4 路独立哈希ahash::RandomState落到每行的一个槽位读取时取所有行计数的最小值作为频率估计。插入与查询均为无锁原子操作fetch_add/loadOrdering::Relaxed时间复杂度 O(h)、空间复杂度 O(h×n)。需要注意计数溢出时可能返回负数调用方需自行处理estimator.rs。因此Rate的observe返回的是估计值而非精确值——这是以极低内存代价换取高吞吐的典型取舍。若需要精确并发计数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::new的Duration与MAX_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),仅供参考
返回列表