资讯动态

Worker Pool模式:高并发任务分发与资源控制

发布时间:2026/9/26 11:01:00 来源:尧图企业网站定制
Worker Pool模式高并发任务分发与资源控制Worker Pool是Go高并发编程中最常用的模式之一。它通过预分配固定数量的goroutine处理任务队列实现资源控制和吞吐量平衡。本文从基础Worker Pool到生产级实现讲透Pool的设计原理、容量控制和优雅关闭。一、核心技术知识点讲解1.1 为什么需要Worker Pool无限制go func()的问题每个goroutine占用2KB栈内存GC压力goroutine越多栈越大GC越慢文件描述符连接/文件操作会耗尽fd资源争抢过多goroutine导致调度开销Worker Pool的优势固定goroutine数资源可控复用goroutine减少分配天然限流任务排队拒绝过载1.2 Pool的基本结构Producer - [Task Channel] - [Worker 1..N] - [Result Channel] - Consumer1.3 动态Pool vs 固定Pool固定Pool启动时创建N个worker适用于任务量稳定的场景动态Pool根据负载动态调整worker数适用于突发流量Go生态ants库提供高性能动态Pool1.4 任务取消与超时Pool必须支持context取消优雅关闭超时控制防止单个任务阻塞整个Pool错误传播收集错误决定是否继续1.5 Pool vs errgrouperrgroup适合并行做N件事Worker Pool适合持续处理M个任务M远大于NPool的生命周期更长适合流式处理1.6 缓冲channel的容量选择任务channel的缓冲大小影响吞吐和延迟缓冲大吞吐高但任务积压严重缓冲小延迟低但producer可能阻塞建议worker数的2-4倍二、实战代码演示2.1 基础Worker Poolpackagemainimport(contextfmtsynctime)typeTaskstruct{IDintDatastring}typeResultstruct{TaskIDintOutputstring}funcworker(ctx context.Context,idint,tasks-chanTask,resultschan-Result,wg*sync.WaitGroup){deferwg.Done()for{select{case-ctx.Done():fmt.Printf(Worker %d: context cancelled\n,id)returncasetask,ok:-tasks:if!ok{fmt.Printf(Worker %d: channel closed, exiting\n,id)return}// 模拟处理result:Result{TaskID:task.ID,Output:fmt.Sprintf(processed(%d): %s,id,task.Data),}results-result}}}funcmain(){ctx,cancel:context.WithTimeout(context.Background(),5*time.Second)defercancel()constnumWorkers3tasks:make(chanTask,10)results:make(chanResult,10)varwg sync.WaitGroupfori:0;inumWorkers;i{wg.Add(1)goworker(ctx,i,tasks,results,wg)}// 派发任务gofunc(){fori:0;i20;i{tasks-Task{ID:i,Data:fmt.Sprintf(task-%d,i)}}close(tasks)}()// 收集结果gofunc(){wg.Wait()close(results)}()forresult:rangeresults{fmt.Printf(Result: %s\n,result.Output)}fmt.Println(All done)}2.2 生产级Pool超时与错误处理packagemainimport(contextfmtsynctime)typeJobstruct{IDintPayloadstringTimeout time.Duration}funcprocessJob(ctx context.Context,job Job)(string,error){ctx,cancel:context.WithTimeout(ctx,job.Timeout)defercancel()select{case-ctx.Done():return,fmt.Errorf(job %d timeout: %w,job.ID,ctx.Err())case-time.After(time.Duration(job.ID%3)*100*time.Millisecond):returnfmt.Sprintf(job-%d-done,job.ID),nil}}typePoolstruct{workersintjobQueuechanJob resultschanstringerrorschanerrorwg sync.WaitGroup}funcNewPool(workers,queueSizeint)*Pool{returnPool{workers:workers,jobQueue:make(chanJob,queueSize),results:make(chanstring,queueSize),errors:make(chanerror,queueSize),}}func(p*Pool)Start(ctx context.Context){fori:0;ip.workers;i{p.wg.Add(1)gop.runWorker(ctx,i)}}func(p*Pool)runWorker(ctx context.Context,idint){deferp.wg.Done()for{select{case-ctx.Done():returncasejob,ok:-p.jobQueue:if!ok{return}result,err:processJob(ctx,job)iferr!nil{p.errors-err}else{p.results-result}}}}func(p*Pool)Submit(job Job){p.jobQueue-job}func(p*Pool)Shutdown(){close(p.jobQueue)p.wg.Wait()close(p.results)close(p.errors)}func(p*Pool)Results()-chanstring{returnp.results}func(p*Pool)Errors()-chanerror{returnp.errors}funcmain(){ctx:context.Background()pool:NewPool(4,100)pool.Start(ctx)gofunc(){fori:0;i30;i{pool.Submit(Job{ID:i,Payload:fmt.Sprintf(data-%d,i),Timeout:2*time.Second,})}pool.Shutdown()}()varerrCountintfor{select{caseresult,ok:-pool.Results():if!ok{fmt.Printf(Done. Errors: %d\n,errCount)return}fmt.Println(OK:,result)caseerr,ok:-pool.Errors():if!ok{continue}fmt.Println(ERR:,err)errCount}}}2.3 使用ants库第三方packagemainimport(fmtsyncsync/atomictime)// ants库使用示例需go get github.com/panjf2000/ants/v2// 这里用简化版本演示typeSimplePoolstruct{taskschanfunc()workersintwg sync.WaitGroup counter atomic.Int64}funcNewSimplePool(workersint)*SimplePool{p:SimplePool{tasks:make(chanfunc{},workers*2),workers:workers,}p.start()returnp}func(p*SimplePool)start(){fori:0;ip.workers;i{p.wg.Add(1)gofunc(){deferp.wg.Done()forfn:rangep.tasks{fn()p.counter.Add(1)}}()}}func(p*SimplePool)Submit(fnfunc()){p.tasks-fn}func(p*SimplePool)Shutdown(){close(p.tasks)p.wg.Wait()}func(p*SimplePool)Completed()int64{returnp.counter.Load()}funcmain(){pool:NewSimplePool(4)varmu sync.Mutex results:make([]int,0,100)fori:0;i100;i{i:i pool.Submit(func(){time.Sleep(10*time.Millisecond)mu.Lock()resultsappend(results,i*2)mu.Unlock()})}pool.Shutdown()fmt.Printf(Completed: %d, Results: %d\n,pool.Completed(),len(results))}2.4 限流Worker Poolpackagemain

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑