Skip to content

Pipeline 流水线模式 ​

难度:⭐⭐⭐ 困难 ​

考点 ​

  • 多阶段 channel 流水线
  • Fan-out / Fan-in
  • 优雅关闭与 goroutine 生命周期管理
  • context 取消传播

题目描述 ​

实现一个数据处理流水线:

  1. Generator — 生成阶段:产生 start 到 end 的整数
  2. Square — 处理阶段:将每个数字平方
  3. Filter — 过滤阶段:只保留满足条件的数字
  4. Merge — 合并阶段:将多个 channel 合并为一个(Fan-in)
  5. 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

提示 ​

  1. 每个阶段启动 goroutine,返回输出 channel
  2. goroutine 中用 select 同时监听 done 和 output channel
  3. Merge 用 sync.WaitGroup 等待所有输入 channel 关闭
  4. 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
}

持续学习,持续构建。