学习目标
学完本章,你应该能够:
- 说清 Kitex 内置限流覆盖哪两个维度(连接数与 QPS),以及它们分别在请求链路的什么时机生效。
- 读懂
ConcurrencyLimiter/RateLimiter/Updatable三个接口的语义,知道Acquire/Release/UpdateLimit怎么用。 - 解释 QPS 限流的"固定窗口 + Token Bucket"算法,能手算给定
limit和interval下每个窗口分多少 token。 - 用
server.WithLimitOption给一个 Kitex 服务配上连接数与 QPS 限流,并理解qpsLimitPostDecode这个开关的取舍。 - 在客户端 / 网关层用
errors.Is(err, kerrors.ErrConnOverLimit)识别限流错误,配合重试做降级。
前置知识:
- Go 基础(接口、goroutine、channel、原子操作概念)。
- Kitex 服务的基本启动方式(
kitex.NewServer+server.WithXxx)。 - 了解"限流"这件事要解决什么问题(保护服务端不被突发流量打垮)。
本章你会动手做的事:
- 给一个 Kitex Server 配置
MaxConnections=10000、MaxQPS=5000,启动后用压测工具打满看拒绝效果。 - 调
UpdateLimit在运行时把 QPS 上限从 5000 改成 2000,验证不重启即生效。 - 用
errors.Is在调用方识别ErrQPSOverLimit,接上重试 / 降级逻辑。
类比:服务端限流就像地铁早高峰的限流栏。站台就那么大,人(连接)太多就塞不进去——
ConcurrencyLimiter管"同时在站里的人数";而检票闸机按秒放行(RateLimiter)管"每秒能过几个人"。两者都是"提前在门口拦住",而不是等挤进车厢(业务 handler)才出问题。
概述
Kitex 提供了内置的服务端限流能力,通过 pkg/limiter 和 pkg/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 --> C3. 动态更新 (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/负载均衡]] — 客户端侧的流量分发
自测题与动手练习
自测题(合上书能答出来,才算懂):
- Kitex 内置限流分哪两个维度?它们分别在请求链路的哪些时机(OnActive / OnRead / OnMessage)生效?
- QPS 限流用了什么算法?当
limit=2000, interval=100ms时,每个 100ms 窗口预分配多少 token? qpsLimitPostDecode取false和true分别意味着什么?默认是哪种、为什么?- 限流被触发时会返回什么错误?调用方怎么识别它?
- 为什么说 Kitex 内置限流"没有分布式能力"?什么场景该换 Sentinel?
动手练习(建议真做一遍):
- 起一个 Kitex 服务,配
MaxConnections=100, MaxQPS=50,用hey或wrk压测,观察超过阈值时客户端收到的错误类型,并区分是连接数还是 QPS 触发的。 - 在管理接口(或测试代码)里调用
UpdateLimit把 QPS 上限动态从 50 改成 20,验证不重启服务即生效,并记录下调前后拒绝率变化。 - 在调用方用
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/超时]](和限流如何配合)。