From 1d53bf2ae1a53e4488eb453aeff0423554771910 Mon Sep 17 00:00:00 2001 From: truewhile <62226914+truewhile@users.noreply.github.com> Date: Wed, 26 Aug 2026 16:16:39 +0800 Subject: [PATCH] 7 --- internal/service/cloud115/queue_executor.go | 12 ++++--- internal/service/strm_queue.go | 13 ++++++-- internal/service/strm_service.go | 37 +++++++++++++++++++++ 3 files changed, 54 insertions(+), 8 deletions(-) diff --git a/internal/service/cloud115/queue_executor.go b/internal/service/cloud115/queue_executor.go index b8e571b..55a7f3e 100644 --- a/internal/service/cloud115/queue_executor.go +++ b/internal/service/cloud115/queue_executor.go @@ -26,13 +26,15 @@ var ( executorOnce sync.Once ) -// GetGlobalExecutor 获取全局队列执行器单例(默认 QPS=8, QPM=480, QPH=20000,保障 115 API 调用安全不超频)。 -// QPS 从 2 提到 8:下载/列表等场景下过去 QPS=2 将所有直链换取串行为每秒 2 个,是下载吞吐的最大瓶颈。 -// 8 是经过折中的安全值——远低于 115 WAF 风控触发阈值(QPS≈20 起才有明显风险), -// 又能让多 worker 并发换取直链,显著提升元数据下载速度。 +// 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(8, 480, 20000) + globalExecutor = NewQueueExecutor(3, 200, 12000) }) return globalExecutor } diff --git a/internal/service/strm_queue.go b/internal/service/strm_queue.go index 8a89fac..f3ff3f3 100644 --- a/internal/service/strm_queue.go +++ b/internal/service/strm_queue.go @@ -28,8 +28,10 @@ const ( ) // downloadWorker 下载队列 worker:认领 → 解析直链 → 下载 → 落盘。 -// 每次批量认领数个任务,并对这批任务并发下载,让「换直链」和「实际下载」 -// 在不同任务间重叠,从而充分利用多线程与 115 换链 QPS。 +// +// 采用「批量认领 + 全局并发限流」:一次认领数个任务,用 StrmService 上的全局信号量 +// 限制整个进程「同时换直链+下载」的并发数(与 115 换链风控匹配,见 strmDownloadSemCap), +// 同时让下载充分并行。换链走全局令牌桶(QPS=3)兜底,下载走 CDN 不限速。 func (s *StrmService) downloadWorker(ctx context.Context) { const claimBatch = 12 // 每次批量认领的任务数 for { @@ -56,12 +58,17 @@ func (s *StrmService) downloadWorker(ctx context.Context) { sleepContext(ctx, 2*time.Second) continue } - // 并发处理本批认领到的任务,充分利用多线程下载 & 直链换取并发 + // 并发处理本批任务:每个任务先获取全局下载槽位,槽位内部执行换链+下载。 + // 信号量与令牌桶双重限速,确保任意时刻并发换链请求不超过安全阈值。 var wg sync.WaitGroup for i := range tasks { wg.Add(1) go func(i int) { defer wg.Done() + if !s.acquireDownloadSlot(ctx) { + return + } + defer s.releaseDownloadSlot() s.processDownloadTask(ctx, &tasks[i]) }(i) } diff --git a/internal/service/strm_service.go b/internal/service/strm_service.go index a43983b..2d46af7 100644 --- a/internal/service/strm_service.go +++ b/internal/service/strm_service.go @@ -90,11 +90,48 @@ type StrmService struct { running map[string]context.CancelFunc // sync path id -> cancel oauthSessions map[string]*strm115AuthSession wafUntil time.Time // 115 风控/限流熔断截止时间(由 mu 保护) + + downloadSem chan struct{} // 全局下载并发信号量:限制整个进程同时进行「换直链+下载」的并发数 + downloadSemOnce sync.Once } // strmWAFCooldown 检测到 115 风控/限流后下载队列的全局冷却时长。 const strmWAFCooldown = 3 * time.Minute +// strmDownloadSemCap 全局同时进行「换直链+下载」的并发上限。 +// +// 115 对换直链接口(/open/ufile/downurl)风控极严:过去把全局 QPS 提到 8 或让多 +// worker 高并发换链,会瞬时撞上 WAF 返回 405 阻断页并触发 180 秒冷却,反而更慢。 +// 因此用信号量把整个进程同时换直链的并发数压到 3,与令牌桶限速共同兜底: +// 宁可下载稍慢,也绝不触发风控。下载本身走 CDN 不限速。 +const strmDownloadSemCap = 3 + +// ensureDownloadSem 惰性初始化全局共享的下载并发信号量。 +func (s *StrmService) ensureDownloadSem() { + s.downloadSemOnce.Do(func() { + s.downloadSem = make(chan struct{}, strmDownloadSemCap) + }) +} + +// acquireDownloadSlot 获取一个下载并发槽位(等待/取消安全)。 +func (s *StrmService) acquireDownloadSlot(ctx context.Context) bool { + s.ensureDownloadSem() + select { + case s.downloadSem <- struct{}{}: + return true + case <-ctx.Done(): + return false + } +} + +// releaseDownloadSlot 释放一个下载并发槽位。 +func (s *StrmService) releaseDownloadSlot() { + if s.downloadSem == nil { + return + } + <-s.downloadSem +} + // NewStrmService constructs the STRM service. func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Container, crypto *CryptoService) *StrmService { return &StrmService{