限速听起来像网关和高并发系统才需要的能力,其实入门项目也经常遇到。比如你要批量同步 5000 个客户资料到第三方 CRM,对方文档写着“每秒最多 5 次请求”。如果你开 50 个 goroutine 一口气打过去,很快就会收到 429,严重时还会被封禁。
Go 标准库没有内置完整的令牌桶限速器,但 time.Ticker 足够写出一个朴素、可理解的小限速器。学习它的过程,也能帮你理解时间驱动、取消和并发控制。
最简单的节拍
每 200ms 执行一次任务,大约就是每秒 5 次:
func syncItems(ctx context.Context, items []Item) error {
ticker := time.NewTicker(200 * time.Millisecond)
defer ticker.Stop()
for _, item := range items {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
if err := syncOne(ctx, item); err != nil {
return err
}
}
}
return nil
}
ticker.Stop() 很重要。Ticker 内部持有资源,不停掉会造成泄漏。虽然小程序马上退出时看不明显,但服务里长期创建 ticker 就会出问题。
第一次是否要等待
上面的代码会在第一个任务前先等 200ms。很多场景希望第一个请求立刻发,后续再按节奏。可以把调用放在等待前,或者先创建一个已经可用的令牌。
func waitTick(ctx context.Context, ticker *time.Ticker) error {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
return nil
}
}
然后在循环里决定什么时候等。简单批处理里,等一下通常无所谓;用户点击按钮后立刻看到响应的场景,第一次延迟就会影响体验。
失败时是否继续
批处理有两种策略:遇到错误立即停止,或者记录错误继续处理下一条。同步第三方数据时,我通常会把失败记录下来,最后统一返回摘要。
type SyncError struct {
ID string
Err error
}
func syncAll(ctx context.Context, items []Item) []SyncError {
ticker := time.NewTicker(200 * time.Millisecond)
defer ticker.Stop()
var errs []SyncError
for _, item := range items {
if err := waitTick(ctx, ticker); err != nil {
errs = append(errs, SyncError{ID: item.ID, Err: err})
break
}
if err := syncOne(ctx, item); err != nil {
errs = append(errs, SyncError{ID: item.ID, Err: err})
}
}
return errs
}
这种写法不会因为一条坏数据让整批停掉,但也不会吞掉错误。最终可以把失败 ID 写入日志、数据库或导出给人工处理。
限速和并发不是一回事
限速控制单位时间内的请求数量,并发控制同时进行的任务数量。只用 ticker 不代表没有并发问题。如果 syncOne 很慢,串行处理会耗时很久;如果你开多个 worker,每个 worker 都一个 ticker,又可能总速率超标。
一种简单方式是所有 worker 共享同一个令牌 channel:
func tokenSource(ctx context.Context, interval time.Duration) <-chan struct{} {
ch := make(chan struct{})
go func() {
defer close(ch)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
select {
case ch <- struct{}{}:
case <-ctx.Done():
return
}
}
}
}()
return ch
}
worker 使用同一个 tokens:
func worker(ctx context.Context, jobs <-chan Item, tokens <-chan struct{}) {
for item := range jobs {
select {
case <-ctx.Done():
return
case <-tokens:
}
_ = syncOne(ctx, item)
}
}
这样不管有多少 worker,总体发起速度都被同一个 token 源控制。
允许短暂突发
有些接口允许“平均每秒 5 次,短时间最多 10 次”。这时可以用带缓冲的 token channel。启动时先放几个令牌,就能允许短暂突发。
func burstTokens(ctx context.Context, interval time.Duration, burst int) <-chan struct{} {
ch := make(chan struct{}, burst)
for i := 0; i < burst; i++ {
ch <- struct{}{}
}
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
close(ch)
return
case <-ticker.C:
select {
case ch <- struct{}{}:
default:
}
}
}
}()
return ch
}
default 表示桶满时丢掉新令牌。这样令牌最多积累到 burst,不会因为系统空闲一小时后突然允许几万次请求。
处理 429
第三方接口返回 429 时,最好尊重 Retry-After:
func retryAfter(resp *http.Response) time.Duration {
v := resp.Header.Get("Retry-After")
if v == "" {
return time.Second
}
if n, err := strconv.Atoi(v); err == nil {
return time.Duration(n) * time.Second
}
if t, err := http.ParseTime(v); err == nil {
return time.Until(t)
}
return time.Second
}
限速器是主动控制,429 是对方告诉你已经超了。两者要一起看。遇到 429 后可以暂停一段时间、降低速率或把任务放回队列,而不是立刻重试。
测试限速逻辑
时间相关代码不好测,因为真实等待会让测试变慢。简单函数可以把 interval 设置很小:
func TestTokenSourceStops(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
tokens := tokenSource(ctx, time.Millisecond)
<-tokens
cancel()
for range tokens {
}
}
更复杂的限速器可以把时钟抽象出来,但入门阶段不必过度设计。先让取消、关闭和基本节奏正确。
小结
time.Ticker 能实现朴素限速:定时产生令牌,任务拿到令牌再执行。它适合批处理、第三方 API 同步、邮件发送和低复杂度后台任务。记得停止 ticker,支持 context 取消,并明确失败策略。
限速不等于并发控制。多个 worker 要共享同一个令牌源,否则总速率会失控。生产项目如果需要精确令牌桶、分布式限速或动态调速,可以再引入成熟库或网关能力。入门阶段先把最小模型写清楚,后面升级才有依据。
常见 Pitfall:Ticker 泄漏
初学者常犯的一个错误是以为 time.NewTicker 创建的 ticker 会在作用域结束时自动回收:
func bad() {
ticker := time.NewTicker(time.Second)
// 没有 defer ticker.Stop()
for i := 0; i < 10; i++ {
<-ticker.C
}
}
ticker 内部在运行时注册了一个定时器。如果不 Stop(),它会一直存在到程序退出。在短时间内反复调用这个函数,内存和定时器都会持续增长。这是 Go 程序里 goroutine 或定时器泄漏的典型例子。
正确做法:
func good() {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
// ...
}
限速策略对比
| 策略 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
time.Ticker | 简单,标准库内置 | 精度有限,不支持动态调速 | 固定的请求间隔 |
| 令牌桶(缓冲 channel) | 支持突发,可共享 | 实现比 ticker 复杂 | API 限流,允许短 burst |
golang.org/x/time/rate | 功能完整,经过验证 | 多一个依赖 | 生产级限速器需求 |
| 信号量 + 计数窗口 | 精确控制并发数 | 实现更复杂 | 并发控制为主的场景 |
对入门项目来说,time.Ticker 加上一个共享 tokens channel 已经覆盖 80% 的场景。只有在需要动态调速、预热、精确令牌桶或分布式限流时才需要引入第三方库。
Ticker 和 Timer 的区别
有人会把 time.Ticker 和 time.Timer 混淆。它们的区别在于:
Timer:只触发一次,触发后 channel 只收到一个值。适合"等待一段时间然后做某事"。Ticker:按固定间隔反复触发。适合"每隔一段时间做某事"。
限速场景必须用 Ticker,Timer 触发一次后就失效了。用错会导致"第一次正确,后续疯狂发请求"的 bug。
// 错误:timer 只触发一次
timer := time.NewTimer(200 * time.Millisecond)
defer timer.Stop()
// 错误:timer 不会重置,循环中一直阻塞
for _, item := range items {
<-timer.C // 只有第一次会触发
}
FAQ:常见限速器问题
Q: 限速器占用多少 CPU?
A: time.Ticker 底层依赖运行时调度,开销很小。即使每秒 1000 次 tick,CPU 占用也几乎不可见。但注意,限速器控制的是"请求发出"的节奏,不是请求本身的执行时间。
Q: 多个 worker 共享 token 时,会不会有争抢问题?
A: channel 的接收是天然互斥的。多个 goroutine 同时从 channel 接收时,只有一个能拿到值。这是 Go 内存模型保证的,不需要额外锁。
Q: 限速和重试怎么配合?
A: 遇到 429 或网络超时后,不要在限速器内部重试。正确的做法是把失败任务丢回队列,让限速器按正常节奏重新调度。否则"限速 + 立刻重试"会让限速形同虚设。
Q: 怎么判断当前限速是否生效?
A: 最直观的方式是在每次请求前后打印日志,看时间戳间隔是否符合预期。更正式的做法是用 Prometheus 等指标暴露请求速率和 429 数量,在 Grafana 上观察趋势。
处理慢请求和积压
限速器保证"每秒发出多少个请求",但不保证请求多快返回。如果下游 API 变慢,发送请求的速度可能还在,但 worker 都在阻塞等待响应。这时可以考虑两层控制:限速 + 超时。
func syncOne(ctx context.Context, item Item) error {
ctx, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, apiURL, body)
resp, err := http.DefaultClient.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
// ...
}
超时可以防止一个慢请求拖住整个 worker。结合限速和超时的批处理程序,通常比单纯限速更稳定。
从日志观察限速效果
限速器是否工作正常,最直观的指标是日志里的时间戳。在每次请求前后打印一条日志:
func syncOne(ctx context.Context, item Item) error {
start := time.Now()
log.Printf("[sync] id=%s start", item.ID)
defer func() {
log.Printf("[sync] id=%s done duration=%v", item.ID, time.Since(start))
}()
// ... 实际请求
}
如果限速器间隔是 200ms,日志里 “start” 的时间戳应该大致按这个间隔分布。如果发现多个请求几乎同时开始,说明限速没有生效。
更正式的做法是用 expvar 或 Prometheus 指标暴露限速相关的计数器,在 Grafana 上观察。
完整限速器示例
把前面学到的东西组合成一个可运行的框架:
type RateLimitedSyncer struct {
client *http.Client
interval time.Duration
}
func NewRateLimitedSyncer(interval time.Duration) *RateLimitedSyncer {
return &RateLimitedSyncer{
client: &http.Client{Timeout: 10 * time.Second},
interval: interval,
}
}
func (s *RateLimitedSyncer) SyncAll(ctx context.Context, items []Item) []SyncError {
tokens := tokenSource(ctx, s.interval)
defer func() {
// 确保 tokenSource goroutine 退出
for range tokens {
}
}()
var errs []SyncError
for _, item := range items {
select {
case <-ctx.Done():
errs = append(errs, SyncError{ID: item.ID, Err: ctx.Err()})
return errs
case <-tokens:
if err := s.syncOne(ctx, item); err != nil {
errs = append(errs, SyncError{ID: item.ID, Err: err})
}
}
}
return errs
}
func (s *RateLimitedSyncer) syncOne(ctx context.Context, item Item) error {
// 实际 HTTP 请求逻辑
return nil
}
这个结构把限速逻辑封装进一个类型,调用方只需要关心业务数据。defer 里的循环消费确保 goroutine 不会因为 channel 还有未读数据而泄漏。
生产环境升级路径
time.Ticker 限速器适合原型和小规模使用。当需求增长时,可以按以下路径升级:
- 需要精确令牌桶:引入
golang.org/x/time/rate - 需要跨进程/分布式限速:使用 Redis + Lua 脚本或专门的限流服务
- 需要按用户/租户分别限速:在 token 源上加一层映射,每个 key 一个限速器
- 需要动态配置:把 interval 和 burst 做成可热更新的配置项
升级时保持接口稳定,内部实现替换即可。这是"先让代码能跑,再逐步优化"的 Go 工程风格。
进阶:自适应限速
有些系统需要"试探性加速"。比如对方文档说每秒 5 次,但实际可以承受 8 次。一种简单自适应策略是:从保守速率开始,如果连续成功就稍微加速,遇到 429 就大幅减速。
type AdaptiveLimiter struct {
mu sync.Mutex
interval time.Duration
minI time.Duration
maxI time.Duration
}
func (a *AdaptiveLimiter) Success() {
a.mu.Lock()
defer a.mu.Unlock()
if a.interval > a.minI {
a.interval = time.Duration(float64(a.interval) * 0.95)
}
}
func (a *AdaptiveLimiter) Throttled() {
a.mu.Lock()
defer a.mu.Unlock()
a.interval = time.Duration(float64(a.interval) * 1.5)
if a.interval > a.maxI {
a.interval = a.maxI
}
}
自适应限速比固定限速更复杂,调试也更难。建议先用固定限速跑稳定,再考虑是否值得升级。
小结
time.Ticker 能实现朴素限速:定时产生令牌,任务拿到令牌再执行。它适合批处理、第三方 API 同步、邮件发送和低复杂度后台任务。记得停止 ticker,支持 context 取消,并明确失败策略。
限速不等于并发控制。多个 worker 要共享同一个令牌源,否则总速率会失控。生产项目如果需要精确令牌桶、分布式限速或动态调速,可以再引入成熟库或网关能力。入门阶段先把最小模型写清楚,后面升级才有依据。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。