并发模式
00:00
Go并发模式详解:工作池、扇出扇入、管道。
1. 工作池
func worker(id int, jobs <-chan int, results chan<- int) {
for j := range jobs {
results <- process(j)
}
}
jobs := make(chan int, 100)
results := make(chan int, 100)
for w := 1; w <= 3; w++ {
go worker(w, jobs, results)
}
for j := 1; j <= 9; j++ {
jobs <- j
}
close(jobs)
2. 扇出扇入
// 扇出:多个 goroutine 处理同一输入
ch1 := fanOut(input)
ch2 := fanOut(input)
ch3 := fanOut(input)
// 扇入:合并多个 channel
merged := fanIn(ch1, ch2, ch3)
func fanIn(channels ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for _, ch := range channels {
wg.Add(1)
go func(c <-chan int) {
for v := range c { out <- v }
wg.Done()
}(ch)
}
go func() { wg.Wait(); close(out) }()
return out
}
3. 管道
func generate(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
}
// 管道组合
ch := square(square(generate(2, 3, 4)))
for v := range ch {
fmt.Println(v) // 16, 81, 256
}