为何所有任务仅在首个goroutine执行?Go扇入扇出管道疑问
Pipeline Fan-In/Fan-Out模式代码执行差异解答
问题描述
我正在实现Pipeline的Fan-In/Fan-Out模式,但不理解两段代码的执行差异,恳请解答:
- 第一段代码:所有任务均在首个goroutine中执行
- 第二段代码:任务大致均匀分配至各个goroutine
第一段代码
func main() { var rdNums []int for i := 0; i < 1000; i++ { rdNums = append(rdNums, i) } pl := generatePipeline(rdNums) c1 := fanOut(pl, "1") c2 := fanOut(pl, "2") c3 := fanOut(pl, "3") c4 := fanOut(pl, "4") c := fanIn(c1, c2, c3, c4) sum := 0 for i := range c { sum += i } fmt.Println(sum) } func generatePipeline(arrNums []int) <-chan int { pl := make(chan int, 100) go func() { for _, n := range arrNums { pl <- n } close(pl) }() return pl } func fanOut(in <-chan int, name string) <-chan int { out := make(chan int) go func() { for v := range in { fmt.Printf("Push square of %d to channel %s \n", v*v, name) out <- v * v } close(out) }() return out } func fanIn(inputChan ...<-chan int) <-chan int { in := make(chan int) go func() { for _, c := range inputChan { for v := range c { in <- v } } close(in) }() return in }
第二段代码
func main() { var rdNums []int for i := 0; i < 1000; i++ { rdNums = append(rdNums, i) } pl := createQueue(rdNums) for i := 0; i < 5; i++ { go process(pl, fmt.Sprintf("worker%d", i)) } time.Sleep(1 * time.Minute) } func createQueue(arrNums []int) <-chan int { pl := make(chan int) go func() { for _, n := range arrNums { pl <- n } close(pl) }() return pl } func process(in <-chan int, name string) { count := 0 go func() { for v := range in { fmt.Printf("Push square of %d to channel %s \n", v*v, name) count++ } fmt.Printf("Process %s success, total number reiceive %d\n", name, count) }() }
差异原因分析
第一段代码:首个goroutine独占任务的核心原因
- Fan-In的串行阻塞逻辑:
fanIn函数中是按顺序遍历每个输出通道,必须完全读完c1的所有数据后,才会开始读取c2、c3、c4的内容。这导致c1对应的goroutine可以持续从输入通道pl中取数据(因为fanIn一直在消费c1的输出,c1的goroutine能不断把计算结果写入自身的out通道,进而可以持续从pl拉取新任务)。 - 输入通道的独占消费:由于
c1的goroutine一直在读取pl,其他fanOut的goroutine根本没有机会从pl中获取任务,最终所有任务都被第一个goroutine执行。即使pl是带缓冲通道,缓冲区内的数据也会被第一个goroutine优先耗尽。
第二段代码:任务均匀分配的关键逻辑
- 多goroutine公平竞争输入通道:多个
process的goroutine同时监听同一个无缓冲通道pl,Go运行时的调度器会采用公平调度策略,将通道中的元素依次分发给不同的goroutine,避免单个goroutine独占所有任务。 - 无额外的串行阻塞:
process函数直接在goroutine中处理从pl获取的数据,没有引入额外的输出通道并被串行读取,每个worker都能独立地从输入通道获取任务,自然实现了任务的均匀分配。
内容的提问来源于stack exchange,提问作者Thang Ha
相关产品推荐
相关产品推荐

