如何在Go的pipeline中使用sync.WaitGroup解决主协程提前退出问题
问题原因
你写的WaitGroup逻辑存在两处核心错误:
wg.Add放在了子goroutine内部执行,可能出现主goroutine已经执行到wg.Wait()时,子goroutine还没调度到wg.Add语句,此时计数器为0,Wait直接返回,主goroutine提前退出wg.Done在子goroutine刚启动就调用了,根本没有等后续的channel读写、业务逻辑执行完成,相当于计数器加1后立刻减1,完全没有起到同步作用
正确实现方案
调整wg.Add的调用位置到启动goroutine之前,并用defer保证wg.Done在子goroutine所有逻辑执行完成后再调用:
package main import ( "fmt" "sync" ) var wg sync.WaitGroup func main() { c1 := make(chan string) c2 := make(chan string) // 启动goroutine前先调用Add,保证计数器在Wait前就完成累加 wg.Add(1) go sender(c1) wg.Add(1) go removeDuplicates(c1, c2) wg.Add(1) go printer(c2) wg.Wait() } func sender(outputStream chan string) { // 用defer保证函数退出前一定调用Done,即使中间出现异常也不会泄漏 defer wg.Done() for _, v := range []string{"one", "one", "two", "two", "three", "three"} { outputStream <- v } close(outputStream) } func removeDuplicates(inputStream, outputStream chan string) { defer wg.Done() temp := "" for v := range inputStream { if v != temp { outputStream <- v temp = v } } close(outputStream) } func printer(inputStream chan string) { defer wg.Done() for v := range inputStream { fmt.Println(v) } }
逻辑说明
调整后三个子goroutine的计数器都会在主goroutine执行Wait前完成累加,每个goroutine执行完全部逻辑(包括channel遍历、数据传递、打印)后才会调用Done减少计数器,主goroutine会等三个计数器全部归0后再退出,完全不需要额外加time.Sleep就能稳定运行。
内容的提问来源于stack exchange,提问作者sv_dev
相关产品推荐
相关产品推荐

