Kitex限流

2024-01-17T14:21:02+08:00 | 8分钟阅读 | 更新于 2024-01-17T14:21:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 说清 Kitex 内置限流覆盖哪两个维度(连接数QPS),以及它们分别在请求链路的什么时机生效。
  2. 读懂 ConcurrencyLimiter / RateLimiter / Updatable 三个接口的语义,知道 Acquire/Release/UpdateLimit 怎么用。
  3. 解释 QPS 限流的"固定窗口 + Token Bucket"算法,能手算给定 limitinterval 下每个窗口分多少 token。
  4. server.WithLimitOption 给一个 Kitex 服务配上连接数与 QPS 限流,并理解 qpsLimitPostDecode 这个开关的取舍。
  5. 在客户端 / 网关层用 errors.Is(err, kerrors.ErrConnOverLimit) 识别限流错误,配合重试做降级。

前置知识

  • Go 基础(接口、goroutine、channel、原子操作概念)。
  • Kitex 服务的基本启动方式(kitex.NewServer + server.WithXxx)。
  • 了解"限流"这件事要解决什么问题(保护服务端不被突发流量打垮)。

本章你会动手做的事

  1. 给一个 Kitex Server 配置 MaxConnections=10000MaxQPS=5000,启动后用压测工具打满看拒绝效果。
  2. UpdateLimit 在运行时把 QPS 上限从 5000 改成 2000,验证不重启即生效。
  3. errors.Is 在调用方识别 ErrQPSOverLimit,接上重试 / 降级逻辑。

类比:服务端限流就像地铁早高峰的限流栏。站台就那么大,人(连接)太多就塞不进去——ConcurrencyLimiter 管"同时在站里的人数";而检票闸机按秒放行(RateLimiter)管"每秒能过几个人"。两者都是"提前在门口拦住",而不是等挤进车厢(业务 handler)才出问题。

概述

Kitex 提供了内置的服务端限流能力,通过 pkg/limiterpkg/limit 包实现,无需依赖外部组件(如 Sentinel)。限流在传输层通过 InboundHandler 拦截器生效,在请求到达业务 handler 之前就进行流量控制。

一句话:它不依赖任何外部中间件,开箱即用、轻量,适合"单实例自我保护"这类诉求。

限流类型

Kitex 支持两种核心限流维度:

类型接口作用拒绝错误
连接数限流ConcurrencyLimiter限制服务端并发连接总数ErrConnOverLimit
QPS 限流RateLimiter限制每秒请求处理速率ErrQPSOverLimit

整体结构长这样:

flowchart LR
    REQ[客户端请求] --> IN[limiterInbound 拦截器]
    IN --> CONN[连接数检查
ConcurrencyLimiter] CONN --> QPS[QPS 检查
RateLimiter] QPS --> H[业务 Handler] QPS -. 超限 .-> E1[返回 ErrQPSOverLimit] CONN -. 超限 .-> E2[返回 ErrConnOverLimit]

1. 连接数限流 (ConcurrencyLimiter)

基于原子计数器实现,跟踪当前活跃连接数:

type ConcurrencyLimiter interface {
    Acquire(ctx context.Context) bool     // 获取连接许可,返回 true 表示允许
    Release(ctx context.Context)          // 释放连接许可
    Status(ctx context.Context) (limit, occupied int)  // 查询当前状态
}
  • Acquire():在连接建立时调用,原子递增计数器,判断是否超过上限
  • Release():在连接关闭时调用,原子递减计数器
  • limit <= 0 时视为不限流模式,始终放行

白话:连接数限流就是"数人头"。每进来一个连接计数器 +1,出去一个 -1,超过上限的新连接直接拒之门外。它保护的是"同时挂着的连接数",防止服务端文件描述符 / 内存被撑爆。

2. QPS 限流 (RateLimiter)

基于固定窗口 + Token Bucket 算法实现:

type RateLimiter interface {
    Acquire(ctx context.Context) bool                          // 获取令牌
    Status(ctx context.Context) (max int, current int, interval time.Duration)  // 查询状态
}

算法原理

  • 将 1 秒划分为多个时间窗口,每个窗口预分配一定数量的 token
  • 使用 time.Ticker 按窗口周期 refill token(上限为总 limit)
  • Acquire() 时原子递减 token 计数,token 耗尽则拒绝
示例:limit=1000, interval=100ms
  每 100ms 窗口分配 1000 / (1000/100) = 100 个 token
  每次 Acquire 消耗 1 个 token
  Ticker 每 100ms 补充一次,上限 1000

用一张图看 token 怎么流动:

flowchart TD
    A[时间按 interval 切片] --> B[每个窗口预分配 token]
    C[time.Ticker 到期] --> D[补充 token 至上限 limit]
    E[请求到达 Acquire] --> F{token 够吗?}
    F -->|够| G[原子递减 放行]
    F -->|不够| H[拒绝 返回 ErrQPSOverLimit]
    G --> C

3. 动态更新 (Updatable)

两个限流器都支持运行时动态调整限流阈值:

type Updatable interface {
    UpdateLimit(limit int)  // 动态修改限流值
}

配合 [[配置管理]] 可实现热更新,无需重启服务。

⚠️ 新手必踩的坑:热更新要小心"瞬间放大"。把 MaxQPS 从 5000 动态调到 20000,下游(DB、缓存)可能根本扛不住这个新流量,等于自己给自己挖坑。热更新一般用于"下游扩容后同步放宽",而不是凭感觉乱调。另外注意:连接数 UpdateLimit 调到比当前已占用还小,新连接会被拒,老连接不受影响,别以为能"立刻缩容"。

服务端集成

基本用法

通过 server.WithLimitOption 在服务启动时配置:

import "github.com/cloudwego/kitex/pkg/limit"

// 步骤 1:构造限流配置,填最大连接数和最大 QPS
svr := kitex.NewServer(YourServiceImpl{},
    server.WithLimitOption(&limit.Option{
        MaxConnections: 10000,  // 最大并发连接数
        MaxQPS:         5000,   // 最大 QPS
    }),
)

拦截器链路

限流通过 pkg/remote/bound.limiterInbound 作为 Inbound Handler 嵌入传输层:

客户端请求
  ↓
OnActive()    → connLimit.Acquire()   // 连接建立时检查连接数
  ↓
OnRead()      → qpsLimit.Acquire()    // 读取请求头后检查 QPS(默认 pre-decode)
  ↓
OnMessage()   → qpsLimit.Acquire()    // 解码完成后检查 QPS(可选 post-decode)
  ↓
业务 Handler  // 通过所有限流检查后进入业务逻辑
  ↓
OnInactive()  → connLimit.Release()   // 连接关闭时释放

把上面这条链路口径画成时序图更直观:

sequenceDiagram
    participant C as 客户端
    participant S as Kitex Server
    participant L as limiterInbound
    participant H as 业务 Handler
    C->>S: 建立连接 OnActive
    S->>L: connLimit.Acquire
    L-->>S: 允许 / 拒绝
    C->>S: 发送请求 OnRead
    S->>L: qpsLimit.Acquire
    L-->>S: 允许 / 拒绝
    S->>H: 进入业务逻辑
    C->>S: 连接关闭 OnInactive
    S->>L: connLimit.Release

关键参数:qpsLimitPostDecode

  • false(默认):在请求解码进行 QPS 限流,保护服务端免受反序列化开销影响
  • true:在请求解码进行 QPS 限流,可根据请求内容做更精细的判断

⚠️ 新手必踩的坑:qpsLimitPostDecode 的默认值是 false(解码前限流)。这意味着超大 / 畸形请求还没解码就被拦了——能保护服务端不被反序列化拖垮,是好事。但如果你想"根据请求体内容决定限流策略"(比如 VIP 用户不限、普通用户限),就得设成 true,代价是要先付出解码开销。按业务需要选,别无脑切。

拒绝处理

当限流触发时,返回标准 Kitex 错误:

// ErrConnOverLimit
//   base: "request over limit"
//   cause: "too many connections"

// ErrQPSOverLimit
//   base: "request over limit"
//   cause: "request too frequent"

可在客户端或网关层通过 errors.Is(err, kerrors.ErrConnOverLimit) 识别限流错误,配合 [[重试]] 机制做降级处理。

架构总结

┌─────────────────────────────────────────────┐
│                 Kitex Server                 │
│                                              │
│  ┌──────────────────────────────────────┐    │
│  │       limiterInbound (Handler)       │    │
│  │  ┌────────────┐  ┌──────────────┐    │    │
│  │  │ConnLimiter │  │  QPSLimiter  │    │    │
│  │  │(原子计数器) │  │(固定窗口+    │    │    │
│  │  │            │  │ TokenBucket) │    │    │
│  │  └─────┬──────┘  └──────┬───────┘    │    │
│  └────────┼────────────────┼────────────┘    │
│           ▼                ▼                  │
│      连接数检查         QPS 检查              │
│           ▼                ▼                  │
│      ┌─────────────────────┐                  │
│      │   Business Handler  │                  │
│      └─────────────────────┘                  │
└─────────────────────────────────────────────┘

与其他方案的对比

方案优点缺点
Kitex 内置限流零依赖、开箱即用、轻量功能较基础,无分布式能力
Apache Sentinel分布式、实时监控、规则中心需要额外部署 Sentinel Dashboard
自定义 Middleware完全可控、灵活需自行实现算法和状态管理

⚠️ 新手必踩的坑:内置限流是"单实例"维度。它只管"这一台 Kitex 进程"的连接数和 QPS,不感知集群里其他实例——也就是说它是本地限流,不是全局分布式限流。如果你的流量入口在网关 / LB 层已经做了全局限流,这里做单实例兜底;如果你期望"全集群共 5000 QPS"那种全局效果,内置限流做不到,得上 Sentinel 之类的分布式方案。

相关笔记

  • [[Kitex/中间件]] — 自定义限流可通过 Middleware 实现
  • [[Kitex/重试]] — 限流拒绝后的客户端重试策略
  • [[Kitex/超时]] — 与限流配合的超时配置
  • [[Kitex/负载均衡]] — 客户端侧的流量分发

自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. Kitex 内置限流分哪两个维度?它们分别在请求链路的哪些时机(OnActive / OnRead / OnMessage)生效?
  2. QPS 限流用了什么算法?当 limit=2000, interval=100ms 时,每个 100ms 窗口预分配多少 token?
  3. qpsLimitPostDecodefalsetrue 分别意味着什么?默认是哪种、为什么?
  4. 限流被触发时会返回什么错误?调用方怎么识别它?
  5. 为什么说 Kitex 内置限流"没有分布式能力"?什么场景该换 Sentinel?

动手练习(建议真做一遍)

  1. 起一个 Kitex 服务,配 MaxConnections=100, MaxQPS=50,用 heywrk 压测,观察超过阈值时客户端收到的错误类型,并区分是连接数还是 QPS 触发的。
  2. 在管理接口(或测试代码)里调用 UpdateLimit 把 QPS 上限动态从 50 改成 20,验证不重启服务即生效,并记录下调前后拒绝率变化。
  3. 在调用方用 errors.Is(err, kerrors.ErrQPSOverLimit) 判断限流错误,接上"退避重试一次 + 仍失败则返回降级结果"的逻辑,跑一遍看效果。

本章小结

  • 两种维度:连接数限流(ConcurrencyLimiter,原子计数器,管"同时在的人头")和 QPS 限流(RateLimiter,固定窗口 + Token Bucket,管"每秒放行速率")。
  • 生效时机早:限流在 limiterInbound 拦截器里、业务 handler 之前就拦,默认在解码前(pre-decode)做 QPS 检查,保护反序列化开销。
  • 可热更新Updatable.UpdateLimit 运行时改阈值,配合配置中心做不重启调整,但放大幅度要谨慎。
  • 识别与降级:拒绝时返回 ErrConnOverLimit / ErrQPSOverLimit,调用方用 errors.Is 识别并接重试 / 降级。
  • 边界清晰:内置限流是单实例本地限流,无分布式能力;需要全局限流请上 Sentinel。

继续可看 [[Kitex/中间件]](自定义限流逻辑)、[[Kitex/重试]](限流后怎么重试降级)、[[Kitex/超时]](和限流如何配合)。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!