You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Go语言滑动窗口限流器测试出现负时间值问题求助

问题描述

我用Golang实现了基于滑动窗口算法的速率限流器,通过map存储每秒的请求数,同时启动独立goroutine清理过期条目。测试时以500毫秒间隔调用AllowRequest方法发起请求,发现输出中出现负的时间差数值,疑惑为何startTime会大于当前时间,请求分析原因。

实现代码

速率限流器接口

type RateLimiterSync interface {
    AllowRequest(userID string) (bool, error)
    Kill() error
}

配置结构体

type Config struct {
    //Duration of timeline
    Duration int
    //no of request allowed
    Count int
}

滑动窗口限流器实现

type SlidingWindowSync struct {
    mutex            sync.Mutex
    keyReqMap        map[string]map[int]int
    keyConfigService *KeyConfigService
    stop             bool
}

func NewSlidingWindowSync(keyConfigService *KeyConfigService, userID string) *SlidingWindowSync {
    var rateLimiter *SlidingWindowSync
    rateLimiter = &SlidingWindowSync{
        keyConfigService: keyConfigService,
        mutex:            sync.Mutex{},
        keyReqMap:        make(map[string]map[int]int),
        stop:             false,
    }
    userConfig, _ := rateLimiter.keyConfigService.GetConfig(userID)
    go rateLimiter.removeEnteries(userID, userConfig.Duration, rateLimiter.mutex)
    return rateLimiter
}

//incrementCountSync : increment req count we can add mutex here
func incrementCountSync(reqMap map[int]int, curTimeInSec int, mutex sync.Mutex) {
    mutex.Lock()
    defer mutex.Unlock()
    if _, exists := reqMap[curTimeInSec]; exists == false {
        reqMap[curTimeInSec] = 1
    } else {
        reqMap[curTimeInSec] = reqMap[curTimeInSec] + 1
    }
}

func (r *SlidingWindowSync) Kill() error {
    var err error
    if r.stop == false {
        r.stop = true
        fmt.Println(" killing go routine")
    } else {
        err = fmt.Errorf(" go routine already stopped")
    }
    return err
}

func (r *SlidingWindowSync) removeEnteries(configKey string, durationInSec int, mutex sync.Mutex) {
    for {
        if r.stop {
            fmt.Println(" existing the go routine ------")
            break
        }
        fmt.Println(" taking a lock in go routine")
        mutex.Lock()
        for key, _ := range r.keyReqMap[configKey] {
            if time.Now().UTC().Second()-key > durationInSec {
                fmt.Println(" deleting key for second =", time.Now().UTC().Second()-key)
                delete(r.keyReqMap[configKey], key)
            }
        }
        mutex.Unlock()
        fmt.Println(" existing a lock in go routine")
        time.Sleep(time.Second * time.Duration(durationInSec) / 10)
    }

}
func getReqCountSync(reqMap map[int]int, curTimeInSec int,
    durationInSec int, mutex sync.Mutex) int {
    var count int
    if len(reqMap) == 0 {
        return 0
    }
    mutex.Lock()
    defer mutex.Unlock()
    //Applying a lock over the critical section of the map
    for key, value := range reqMap {
        if curTimeInSec-key < durationInSec {
            count += value
        }
    }
    return count
}

//AllowRequest : this is the main impl for the rate limiter
func (s *SlidingWindowSync) AllowRequest(userID string) (bool, error) {
    var (
        reqMap     map[int]int
        userConfig *domains.Config
        err        error
        curTime    = time.Now().UTC().Second()
    )
    //Get the config service
    if userConfig, err = s.keyConfigService.GetConfig(userID); err != nil {
        return false, err
    }
    //Check in map for req count
    if _, exists := s.keyReqMap[userID]; !exists {
        s.keyReqMap[userID] = make(map[int]int)
    }
    reqMap = s.keyReqMap[userID]
    if getReqCountSync(reqMap, curTime, userConfig.Duration, s.mutex) >= userConfig.Count {
        return false, nil
    }
    incrementCountSync(reqMap, curTime, s.mutex)
    return true, nil
}

测试代码

func main() {
    keyConfigService := services.NewKeyConfigService()
    keyConfigService.Upsert("login", domains.NewConfig(10, 10))
    rateLimiter := services.NewSlidingWindowSync(keyConfigService, "login")
    startTime := time.Now().UTC().Second()
    for index := 0; index < 40; index++ {
        if allow, err := rateLimiter.AllowRequest("login"); err != nil {
            fmt.Errorf(" error = %+v", err)
        } else if allow {
            fmt.Printf(" request allowed at %d second \n",
                time.Now().UTC().Second()-startTime)
        } else {
            fmt.Printf(" request aborted at %d second \n",
                time.Now().UTC().Second()-startTime)
        }
        time.Sleep(time.Millisecond * 500)
    }

}

测试输出

request aborted at -44 second 
request aborted at -44 second 
...
原因分析与解决方法

核心问题:时间取值错误

你用time.Now().UTC().Second()获取的是当前时间分钟内的秒数(0-59),不是从Unix纪元开始的总秒数。当测试跨分钟边界时,比如startTime是某分钟的59秒,过1秒后进入下一分钟,此时time.Now().UTC().Second()返回0,计算0 - 59就会得到负数,这就是输出负时间差的原因。

额外问题:Mutex值传递失效

代码中incrementCountSync、removeEnteries、getReqCountSync函数接收的是sync.Mutex值类型参数,这会导致每次调用都会复制一个新的Mutex,多个goroutine实际使用的是不同的锁,完全起不到同步保护作用,会引发数据竞争问题。

修复方案

  1. 替换时间获取方式:把所有time.Now().UTC().Second()替换为time.Now().UTC().Unix(),这个方法返回从Unix纪元开始的总秒数,是持续递增的,不会出现循环的0-59,也就不会产生负数时间差。
  2. 修正Mutex传递方式:将所有接收Mutex的函数参数改为*sync.Mutex指针类型,确保所有goroutine使用同一个锁,正确保护共享的map数据。

示例修正后的AllowRequest中时间部分:

curTime = time.Now().UTC().Unix()

修正后的incrementCountSync函数:

func incrementCountSync(reqMap map[int64]int, curTimeInSec int64, mutex *sync.Mutex) {
    mutex.Lock()
    defer mutex.Unlock()
    if _, exists := reqMap[curTimeInSec]; exists == false {
        reqMap[curTimeInSec] = 1
    } else {
        reqMap[curTimeInSec] = reqMap[curTimeInSec] + 1
    }
}

注意:同时需要把相关的int类型时间参数和map的key类型改为int64,因为Unix()返回的是int64类型。

内容的提问来源于stack exchange,提问作者nitin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 21:53:25