You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.05 15:18:03