Pipeline 流水线模式
难度:⭐⭐⭐ 困难
考点
- 多阶段 channel 流水线
- Fan-out / Fan-in
- 优雅关闭与 goroutine 生命周期管理
- context 取消传播
题目描述
实现一个数据处理流水线:
Generator— 生成阶段:产生 start 到 end 的整数Square— 处理阶段:将每个数字平方Filter— 过滤阶段:只保留满足条件的数字Merge— 合并阶段:将多个 channel 合并为一个(Fan-in)Pipeline— 组合以上阶段,构建完整流水线
所有阶段必须:
- 接受 done channel 用于取消
- 在 done 关闭后立即退出,不泄漏 goroutine
函数签名
go
func Generator(done <-chan struct, start, end int) <-chan int
func Square(done <-chan struct{}, in <-chan int) <-chan int
func Filter(done <-chan struct{}, in <-chan int, predicate func(int) bool) <-chan int
func Merge(done <-chan struct{}, channels ...<-chan int) <-chan int
func Pipeline(start, end int, predicate func(int) bool) []int提示
- 每个阶段启动 goroutine,返回输出 channel
- goroutine 中用 select 同时监听 done 和 output channel
- Merge 用 sync.WaitGroup 等待所有输入 channel 关闭
- Pipeline 创建 done channel,可以通过 close(done) 取消整条流水线
参考答案(Go)
点击展开参考答案
go
//go:build ignore
package answer
import "sync"
func Generator(done <-chan struct{}, start, end int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for i := start; i <= end; i++ {
select {
case out <- i:
case <-done:
return
}
}
}()
return out
}
func Square(done <-chan struct{}, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for v := range in {
select {
case out <- v * v:
case <-done:
return
}
}
}()
return out
}
func Filter(done <-chan struct{}, in <-chan int, predicate func(int) bool) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for v := range in {
if predicate(v) {
select {
case out <- v:
case <-done:
return
}
}
}
}()
return out
}
func Merge(done <-chan struct{}, 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) {
defer wg.Done()
for v := range c {
select {
case out <- v:
case <-done:
return
}
}
}(ch)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func Pipeline(start, end int, predicate func(int) bool) []int {
done := make(chan struct{})
defer close(done)
gen := Generator(done, start, end)
squared := Square(done, gen)
filtered := Filter(done, squared, predicate)
var result []int
for v := range filtered {
result = append(result, v)
}
if result == nil {
return []int{}
}
return result
}