Files
MeBox/internal/service/cloud115/queue_executor.go
T
truewhile 1d53bf2ae1 7
2026-08-26 16:16:39 +08:00

112 lines
2.9 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package cloud115
import (
"context"
"sync"
"time"
"golang.org/x/time/rate"
)
// QueueExecutor 全局 115 请求调度执行器,通过三级令牌桶(QPS / QPM / QPH)平滑所有外发请求。
type QueueExecutor struct {
sync.RWMutex
qpsLimiter *rate.Limiter
qpmLimiter *rate.Limiter
qphLimiter *rate.Limiter
throttleManager *ThrottleManager
qpsConfig int
qpmConfig int
qphConfig int
}
var (
globalExecutor *QueueExecutor
executorOnce sync.Once
)
// GetGlobalExecutor 获取全局队列执行器单例(默认 QPS=3, QPM=200, QPH=12000,保障 115 API 调用安全不超频)。
//
// 历史教训:QPS 提到 8 后,下载换直链接口(/open/ufile/downurl,WAF 重点盯防对象)
// 瞬时突发撞上 115 风控,返回阿里云 405 阻断页(HTTP 405),导致全量同步失败。
// 因此回调到 3——这是经过实测的安全上限:宁慢勿触发风控,一旦 405 冷却 180 秒,
// 整体吞吐反而更低。下载实际走 CDN 不受此限速影响,瓶颈仅在换链环节。
func GetGlobalExecutor() *QueueExecutor {
executorOnce.Do(func() {
globalExecutor = NewQueueExecutor(3, 200, 12000)
})
return globalExecutor
}
// NewQueueExecutor 创建新的队列调度执行器。
func NewQueueExecutor(qps, qpm, qph int) *QueueExecutor {
if qps <= 0 {
qps = 3
}
if qpm <= 0 {
qpm = 200
}
if qph <= 0 {
qph = 12000
}
return &QueueExecutor{
qpsLimiter: rate.NewLimiter(rate.Limit(qps), qps),
qpmLimiter: rate.NewLimiter(rate.Every(time.Minute/time.Duration(qpm)), qpm),
qphLimiter: rate.NewLimiter(rate.Every(time.Hour/time.Duration(qph)), qph),
throttleManager: GetGlobalThrottleManager(),
qpsConfig: qps,
qpmConfig: qpm,
qphConfig: qph,
}
}
// Acquire 统一在发送请求前获取令牌,并等待熔断恢复。
func (qe *QueueExecutor) Acquire(ctx context.Context) error {
// 1. 如果处于限流状态,先阻塞等待 60s 静默冷却结束
if err := qe.throttleManager.WaitThrottleRecovery(ctx); err != nil {
return err
}
// 2. 依次获取小时级、分钟级、秒级令牌
if err := qe.qphLimiter.Wait(ctx); err != nil {
return err
}
if err := qe.qpmLimiter.Wait(ctx); err != nil {
return err
}
if err := qe.qpsLimiter.Wait(ctx); err != nil {
return err
}
return nil
}
// SetRateLimitConfig 动态调整速率限制。
func (qe *QueueExecutor) SetRateLimitConfig(qps, qpm, qph int) {
qe.Lock()
defer qe.Unlock()
if qps > 0 {
qe.qpsConfig = qps
qe.qpsLimiter.SetLimit(rate.Limit(qps))
qe.qpsLimiter.SetBurst(qps)
}
if qpm > 0 {
qe.qpmConfig = qpm
qe.qpmLimiter.SetLimit(rate.Every(time.Minute / time.Duration(qpm)))
qe.qpmLimiter.SetBurst(qpm)
}
if qph > 0 {
qe.qphConfig = qph
qe.qphLimiter.SetLimit(rate.Every(time.Hour / time.Duration(qph)))
qe.qphLimiter.SetBurst(qph)
}
}
// MarkThrottled 触发限流熔断。
func (qe *QueueExecutor) MarkThrottled() {
qe.throttleManager.MarkThrottled()
}