// Package rate 提供按秒配额的带宽限速,常配合 conn 包使用。 package rate import ( "sync/atomic" "time" ) // Rate 带宽限速器:每秒重置配额,配额用尽时 Add 阻塞等待 type Rate struct { Limit int32 `json:"limit"` // 限制速度 字节 0为不限制 Surplus int32 `json:"surplus"` // 剩余流量 字节 Now int32 `json:"now"` // 接收(下载/下行) 当前流量 字节 Max int32 `json:"max"` // 接收(下载/下行) 最大流量 字节 IsStop chan bool `json:"-"` } // NewRate 创建限速器,bandwidth 为带宽上限(单位 Mbps) func NewRate(bandwidth int) *Rate { // 带宽 应该是Mbps 要换算成Mbyte limit := bandwidth * 1024 * 1024 / 8 r := new(Rate) r.Limit = int32(limit) r.Now = 0 r.Surplus = int32(limit) r.IsStop = make(chan bool) // 初始化停止通道,否则 Stop() 会 panic return r } // SetLimit 动态调整带宽上限,bandwidth 为 Mbps(0 表示不限速) func (r *Rate) SetLimit(bandwidth int) { // 带宽 应该是Mbps 要换算成Mbyte limit := 0 if bandwidth > 0 { limit = bandwidth * 1024 * 1024 / 8 } else { limit = 0 } atomic.StoreInt32(&r.Limit, int32(limit)) } // Start 启动每秒配额重置协程 func (r *Rate) Start() { go r.proc() } // Stop 停止配额重置协程 func (r *Rate) Stop() { r.IsStop <- true } // proc 配额重置循环 func (r *Rate) proc() { ticker := time.NewTicker(time.Second * 1) for { select { case <-ticker.C: r.reset() case <-r.IsStop: ticker.Stop() return } } } // reset 统计上一秒用量并重置配额 func (r *Rate) reset() { n := r.Limit - atomic.LoadInt32(&r.Surplus) atomic.StoreInt32(&r.Now, n) atomic.StoreInt32(&r.Surplus, r.Limit) if n > 0 { now := atomic.LoadInt32(&r.Now) if now > 0 && now > atomic.LoadInt32(&r.Max) { atomic.StoreInt32(&r.Max, now) } } //fmt.Printf("now:%d limit:%d Surplus:%d \n", n, r.Limit, r.Limit) } // GetNow 获取当前秒已用量(字节) func (r *Rate) GetNow() int { n := atomic.LoadInt32(&r.Now) return int(n) } // ResetMax 取出并清零当前秒用量 func (r *Rate) ResetMax() int { n := atomic.SwapInt32(&r.Now, 0) return int(n) } // Add 计入本秒用量;配额用尽时阻塞等待下一秒配额释放 func (r *Rate) Add(size int) { if atomic.LoadInt32(&r.Surplus) > 0 || atomic.LoadInt32(&r.Limit) == 0 { atomic.AddInt32(&r.Surplus, -int32(size)) return } for { //fmt.Println("等待") time.Sleep(time.Millisecond * 10) if atomic.LoadInt32(&r.Surplus) > 0 { atomic.AddInt32(&r.Surplus, -int32(size)) return } } }