0Pricing
Go Academy · 课时

扇出与扇入

分发并收集工作

扇出与扇入 是 CoddyKit 上的免费 Go Academy 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Go Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Go Academy 课程共包含 4 节课。

扇出与扇入

扇出是指启动多个 Go 协程,从一个通道读取数据并分发工作。扇入是指将多个 Go 协程的输出重新合并到一个通道中。

  • 扇出会将负载分散到多个 worker
  • 扇入会为使用者汇总结果

生成器 stage

管道从生成器开始:这是一个返回通道的函数,它在自己的 Go 协程中向通道写入值,然后关闭该通道。

func gen(nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        for _, n := range nums {
            out <- n
        }
        close(out)
    }()
    return out
}

执行扇出

要执行扇出,请启动多个从同一个输入通道读取数据的 Go 协程。每个协程都运行相同的处理函数,并生成自己的输出通道。

in := gen(1, 2, 3, 4, 5)
c1 := square(in)
c2 := square(in)

处理 stage

square stage 从 in 读取数据,将每个值平方,然后发送到新的输出通道。此 stage 的多个实例共享一个输入通道。

func square(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        for n := range in {
            out <- n * n
        }
        close(out)
    }()
    return out
}

使用 merge 执行扇入

扇入会将多个通道合并为一个通道。sync.WaitGroup 会跟踪每个通道对应的 Go 协程,而关闭协程会等待所有协程完成后,再关闭合并后的输出通道。

func merge(cs ...<-chan int) <-chan int {
    out := make(chan int)
    var wg sync.WaitGroup
    for _, c := range cs {
        wg.Add(1)
        go func(c <-chan int) {
            for n := range c {
                out <- n
            }
            wg.Done()
        }(c)
    }
    go func() { wg.Wait(); close(out) }()
    return out
}

为什么需要 WaitGroup

每个负责合并的 Go 协程都会将一个输入通道复制到 out。在所有输入通道都读取完之前,不能关闭 out。WaitGroup 会统计活动中的复制协程;关闭协程会阻塞在 wg.Wait() 上,直到这些协程完成。

捕获循环变量

在较早版本的 Go 中,闭包必须将 c 作为参数接收,这样每个 Go 协程绑定的就是自己的通道,而不是共享的循环变量。显式传递 c 可以避免经典的循环变量错误。

完整的扇出与扇入

此程序生成 numbers,分发到两个 square stage,然后将结果扇入并汇总合并后的结果。

package main

import (
    "fmt"
    "sync"
)

func gen(nums ...int) <-chan int {
    out := make(chan int)
    go func() { for _, n := range nums { out <- n }; close(out) }()
    return out
}

func square(in <-chan int) <-chan int {
    out := make(chan int)
    go func() { for n := range in { out <- n * n }; close(out) }()
    return out
}

func merge(cs ...<-chan int) <-chan int {
    out := make(chan int)
    var wg sync.WaitGroup
    for _, c := range cs {
        wg.Add(1)
        go func(c <-chan int) { for n := range c { out <- n }; wg.Done() }(c)
    }
    go func() { wg.Wait(); close(out) }()
    return out
}

func main() {
    in := gen(1, 2, 3, 4, 5)
    c1 := square(in)
    c2 := square(in)
    total := 0
    for n := range merge(c1, c2) {
        total += n
    }
    fmt.Println("sum of squares:", total)
}

顺序并不保证

由于两个 square stage 会竞争读取同一个输入,合并后的输出顺序是不确定的。如果需要保持顺序,请为每个项目附加索引,或使用单个 stage。

何时使用扇出

当某个 stage 成为瓶颈且工作可以并行处理时,可以使用扇出。如果上游生成器速度较慢,增加更多下游 worker 并没有帮助,它们只会处于饥饿状态。

背压

无缓冲通道会产生自然的背压:当速度较慢的使用者无法跟上时,速度较快的 stage 就会阻塞。这可以防止整个管道中的内存无限增长。

快速检查

请测试您对扇入的理解。

回顾

您已学会扇出与扇入:

  • 扇出:多个 Go 协程读取同一个通道
  • 扇入:通过 WaitGroup 将多个通道合并为一个
  • 并行 stage 不会保留原有顺序
  • 无缓冲通道提供背压

常见问题解答

「扇出与扇入」课时是免费的吗?

是的 — 「扇出与扇入」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Go Academy 课程的其余内容,请升级到 CoddyKit PRO。 Go Academy 课程共包含 4 节课。

「扇出与扇入」这节课中我会学到什么?

分发并收集工作 你通过在浏览器中直接运行的动手代码来练习 Go Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Go Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Go Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。

「扇出与扇入」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Go Academy 课中编写并运行代码吗?

能。每节 Go Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 工作池模式
  2. 扇出与扇入
  3. 处理流程阶段
  4. 优雅关闭
← 返回 Go Academy