WaitGroup与errgroup:并发任务编排
发布时间:2026/9/26 11:00:50来源:尧图网络
WaitGroup与errgroup并发任务编排sync.WaitGroup是并发编程的标配errgroup在其之上增加了错误传播和取消语义。从WaitGroup的计数器语义Add/Done/Wait到errgroup的WithContext取消本文讲透并发任务编排的底层机制与最佳实践。一、核心技术知识点讲解1.1 WaitGroup的底层结构typeWaitGroupstruct{noCopy noCopy state atomic.Uint64// 高32位计数器低32位等待者数semauint32// 信号量}关键点计数器高32位Add/Done增减等待者数低32位Wait的goroutine数信号量Wait阻塞和唤醒机制1.2 Add/Done/Wait语义varwg sync.WaitGroup wg.Add(1)// 计数器1gofunc(){deferwg.Done()// 计数器-1}()wg.Wait()// 阻塞直到计数器为0Add必须在Wait之前或并发安全地调用计数器为正时Done等价于Add(-1)计数器不能为负panic: negative WaitGroup counter1.3 WaitGroup的陷阱Add与Wait的竞争不能在Wait的同时Add未使用race检测可能panicDone次数过多panic复制WaitGroup禁止复制noCopy标记go vet检测Wait之后重用需确保计数器归零后才能安全重用1.4 errgroup.Groupimportgolang.org/x/sync/errgroupvarg errgroup.Group g.Go(func()error{// 任务returnnil})iferr:g.Wait();err!nil{// 处理错误}第一个返回错误的goroutine会取消其他任务WithContext自动传播第一个非nil错误限制并发数g.SetLimit(n)1.5 errgroup.WithContextg,ctx:errgroup.WithContext(ctx)g.Go(func()error{select{case-ctx.Done():returnctx.Err()// 其他goroutine出错取消default:}// 业务returnnil})1.6 编排模式对比工具用途错误传播取消WaitGroup简单等待无无errgroup错误传播取消第一个错误WithContexterrgroupSetLimit并发限制第一个错误支持channel自定义编排手动手动二、实战代码演示2.1 WaitGroup基础packagemainimport(fmtsynctime)funcmain(){varwg sync.WaitGroupfori:0;i5;i{wg.Add(1)// 在主goroutine中Addgofunc(iint){deferwg.Done()time.Sleep(time.Duration(i)*100*time.Millisecond)fmt.Printf(Task %d done\n,i)}(i)}wg.Wait()fmt.Println(All tasks done)}2.2 WaitGroup的Add时机packagemainimport(fmtsynctime)funccorrectAdd(){varwg sync.WaitGroup jobs:[]int{1,2,3,4,5}wg.Add(len(jobs))// 统一Addfor_,j:rangejobs{gofunc(jobint){deferwg.Done()time.Sleep(10*time.Millisecond)fmt.Printf(Processed %d\n,job)}(j)}wg.Wait()fmt.Println(Done)}funcmain(){correctAdd()}2.3 errgroup基础packagemainimport(errorsfmttimegolang.org/x/sync/errgroup)funcmain(){varg errgroup.Group g.Go(func()error{time.Sleep(100*time.Millisecond)fmt.Println(Task 1 done)returnnil})g.Go(func()error{time.Sleep(200*time.Millisecond)fmt.Println(Task 2 done)returnnil})g.Go(func()error{time.Sleep(50*time.Millisecond)returnerrors.New(task 3 failed)})err:g.Wait()iferr!nil{fmt.Printf(Error: %v\n,err)}else{fmt.Println(All succeeded)}}2.4 errgroup.WithContext取消packagemainimport(contexterrorsfmttimegolang.org/x/sync/errgroup)funcmain(){g,ctx:errgroup.WithContext(context.Background())g.Go(func()error{select{case-ctx.Done():fmt.Println(Task 1 cancelled:,ctx.Err())returnctx.Err()case-time.After(500*time.Millisecond):fmt.Println(Task 1 completed)returnnil}})g.Go(func()error{time.Sleep(100*time.Millisecond)returnerrors.New(task 2 failed fast)})g.Go(func()error{select{case-ctx.Done():fmt.Println(Task 3 cancelled:,ctx.Err())returnctx.Err()case-time.After(3*time.Second):fmt.Println(Task 3 completed)returnnil}})err:g.Wait()fmt.Printf(Final error: %v\n,err)}2.5 生产级并行API聚合packagemainimport(contextfmttimegolang.org/x/sync/errgroup)typeProfilestruct{UserIDintNamestringOrders[]stringFollowersint}funcfetchUserName(ctx context.Context,idint)(string,error){time.Sleep(50*time.Millisecond)returnfmt.Sprintf(user-%d,id),nil}funcfetchOrders(ctx context.Context,idint)([]string,error){time.Sleep(80*time.Millisecond)return[]string{order-1,order-2},nil}funcfetchFollowers(ctx context.Context,idint)(int,error){time.Sleep(30*time.Millisecond)return1000,nil}funcgetProfileParallel(ctx context.Context,userIDint)(*Profile,error){g,ctx:errgroup.WithContext(ctx)profile:Profile{UserID:userID}g.Go(func()error{name,err:fetchUserName(ctx,userID)iferr!nil{returnerr}profile.Namenamereturnnil})g.Go(func()error{orders,err:fetchOrders(ctx,userID)iferr!nil{returnerr}profile.Ordersordersreturnnil})g.Go(func()error{followers,err:fetchFollowers(ctx,userID)iferr!nil{returnerr}profile.Followersfollowersreturnnil})iferr:g.Wait();err!nil{returnnil,err}returnprofile,nil}funcmain(){start:time.Now()p,err:getProfileParallel(context.Background(),42)iferr!nil{fmt.Println(Error:,err)return}fmt.Printf(Profile: %v (took %v)\n,p,time.Since(start))}2.6 errgroup并发限制packagemainimport(contextfmtsynctimegolang.org/x/sync/errgroup)funcmain(){g,_:errgroup.WithContext(context.Background())g.SetLimit(3)varmu sync.Mutexvaractiveintfori:0;i10;i{i:i g.Go(func()error{mu.Lock()activecurrent:active mu.Unlock()fmt.Printf(Task %d starting (active: %d)\n,i,current)time.Sleep(200*time.Millisecond)mu.Lock()active--mu.Unlock()fmt.Printf(Task %d done\n,i)returnnil})}iferr:g.Wait();err!nil{fmt.Println(Error:,err)}fmt.Println(All done)}三、开发痛点与报错避坑指南3.1 复制WaitGroupfuncprocess(wg sync.WaitGroup){// ❌ 值传递复制deferwg.Done()}避坑方案用指针传递funcprocess(wg*sync.WaitGroup){deferwg.Done()}3.2 negative WaitGroup countervarwg sync.WaitGroup wg.Add(1)wg.Done()wg.Done()// panic: negative WaitGroup counter避坑方案每个Add对应恰好一个Done用defer保证。3.3 Wait与Add竞争痛点描述goroutine内Add主goroutine Wait可能发生Wait先执行完计数器为0然后Add导致goroutine没被等待。避坑方案所有Add在主goroutine中完成或Add在goroutine启动前完成。3.4 errgroup的取消语义误解痛点描述认为g.Wait返回错误后所有任务立即停止。实际上只有使用ctx.Done()检查的任务才会响应取消。避坑方案任务中监听ctx.Done()纯计算任务不会被打断。3.5 SetLimit与死锁痛点描述SetLimit(n)太小且任务内再启动g.Go可能死锁等待goroutine占满限制。避坑方案内部任务用独立group或限制1且任务不自旋。四、全文总结WaitGroup结构计数器等待者数信号量atomic实现。使用准则Add在主goroutine、Done用defer、禁止复制。errgroup增强错误传播WithContext取消SetLimit并发限制。取消机制第一个错误触发ctx取消任务需监听ctx.Done()。应用场景并行API聚合、批量任务、优雅关闭。对比选择简单等待用WaitGroup需要错误/取消用errgroup。五、技术进阶展望泛型errgroup未来可能提供泛型版本的返回值聚合。结构化并发conc库的WaitGroup增强panic捕获、结果收集。限流集成errgroup.SetLimit与令牌桶结合的并发控制。编排模式库flow、sema等库提供更丰富的编排原语。六、参考文献Go官方文档 - sync.WaitGroup: https://pkg.go.dev/sync#WaitGroupGo源码 - sync/waitgroup.go: https://github.com/golang/go/blob/master/src/sync/waitgroup.goerrgroup文档: https://pkg.go.dev/golang.org/x/sync/errgroupGo Blog - Pipelines and cancellation: https://go.dev/blog/pipelinessourcegraph/conc: https://github.com/sourcegraph/conc
网站建设高端定制企业官网