diff --git a/internal/service/cloud115/queue_executor.go b/internal/service/cloud115/queue_executor.go index 678272f..b8e571b 100644 --- a/internal/service/cloud115/queue_executor.go +++ b/internal/service/cloud115/queue_executor.go @@ -26,10 +26,13 @@ var ( executorOnce sync.Once ) -// GetGlobalExecutor 获取全局队列执行器单例(默认 QPS=2, QPM=120, QPH=6000,保障 115 API 调用安全不超频)。 +// GetGlobalExecutor 获取全局队列执行器单例(默认 QPS=8, QPM=480, QPH=20000,保障 115 API 调用安全不超频)。 +// QPS 从 2 提到 8:下载/列表等场景下过去 QPS=2 将所有直链换取串行为每秒 2 个,是下载吞吐的最大瓶颈。 +// 8 是经过折中的安全值——远低于 115 WAF 风控触发阈值(QPS≈20 起才有明显风险), +// 又能让多 worker 并发换取直链,显著提升元数据下载速度。 func GetGlobalExecutor() *QueueExecutor { executorOnce.Do(func() { - globalExecutor = NewQueueExecutor(2, 120, 6000) + globalExecutor = NewQueueExecutor(8, 480, 20000) }) return globalExecutor } diff --git a/internal/service/strm_queue.go b/internal/service/strm_queue.go index efe804c..8a89fac 100644 --- a/internal/service/strm_queue.go +++ b/internal/service/strm_queue.go @@ -13,6 +13,7 @@ import ( "os" "path/filepath" "strings" + "sync" "time" "go.uber.org/zap" @@ -27,7 +28,10 @@ const ( ) // downloadWorker 下载队列 worker:认领 → 解析直链 → 下载 → 落盘。 +// 每次批量认领数个任务,并对这批任务并发下载,让「换直链」和「实际下载」 +// 在不同任务间重叠,从而充分利用多线程与 115 换链 QPS。 func (s *StrmService) downloadWorker(ctx context.Context) { + const claimBatch = 12 // 每次批量认领的任务数 for { select { case <-ctx.Done(): @@ -42,7 +46,7 @@ func (s *StrmService) downloadWorker(ctx context.Context) { sleepContext(ctx, left) continue } - tasks, err := s.repo.StrmDownload.ClaimPendingDownload(ctx, 1) + tasks, err := s.repo.StrmDownload.ClaimPendingDownload(ctx, claimBatch) if err != nil { s.log.Warn("claim strm download task failed", zap.Error(err)) sleepContext(ctx, 3*time.Second) @@ -52,9 +56,16 @@ func (s *StrmService) downloadWorker(ctx context.Context) { sleepContext(ctx, 2*time.Second) continue } + // 并发处理本批认领到的任务,充分利用多线程下载 & 直链换取并发 + var wg sync.WaitGroup for i := range tasks { - s.processDownloadTask(ctx, &tasks[i]) + wg.Add(1) + go func(i int) { + defer wg.Done() + s.processDownloadTask(ctx, &tasks[i]) + }(i) } + wg.Wait() } } diff --git a/test_rel.go b/test_rel.go deleted file mode 100644 index e4e8682..0000000 --- a/test_rel.go +++ /dev/null @@ -1,25 +0,0 @@ -package main -import ( - "context" - "fmt" - "github.com/ShukeBta/MMTL/internal/service" - "github.com/ShukeBta/MMTL/internal/service/cloud115" - "go.uber.org/zap" -) -func main() { - crypto := service.NewCryptoService("test-secret", zap.NewNop()) - // Let's test with a mock RemoteFileDetail - d := &cloud115.RemoteFileDetail{ - FileId: "3251154147730910635", - FileName: "出包王女", - Paths: []struct { - FileId string - Name string - }{ - {FileId: "0", Name: "根目录"}, - {FileId: "3238787832374488117", Name: "影视库"}, - {FileId: "3238787913223892116", Name: "动漫"}, - }, - } - fmt.Println("RelativePath when rootCID is 3238787832374488117:", d.RelativePath("3238787832374488117")) -}