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

为何所有任务仅在首个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独占任务的核心原因

  1. Fan-In的串行阻塞逻辑:fanIn函数中是按顺序遍历每个输出通道,必须完全读完c1的所有数据后,才会开始读取c2、c3、c4的内容。这导致c1对应的goroutine可以持续从输入通道pl中取数据(因为fanIn一直在消费c1的输出,c1的goroutine能不断把计算结果写入自身的out通道,进而可以持续从pl拉取新任务)。
  2. 输入通道的独占消费:由于c1的goroutine一直在读取pl,其他fanOut的goroutine根本没有机会从pl中获取任务,最终所有任务都被第一个goroutine执行。即使pl是带缓冲通道,缓冲区内的数据也会被第一个goroutine优先耗尽。

第二段代码:任务均匀分配的关键逻辑

  1. 多goroutine公平竞争输入通道:多个process的goroutine同时监听同一个无缓冲通道pl,Go运行时的调度器会采用公平调度策略,将通道中的元素依次分发给不同的goroutine,避免单个goroutine独占所有任务。
  2. 无额外的串行阻塞:process函数直接在goroutine中处理从pl获取的数据,没有引入额外的输出通道并被串行读取,每个worker都能独立地从输入通道获取任务,自然实现了任务的均匀分配。

内容的提问来源于stack exchange,提问作者Thang Ha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 09:14:56