扇出与扇入
分发并收集工作
扇出与扇入 是 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 反馈 — 无需本地设置。