go-Channel并发模式-管道 Pipeline

管道:多个阶段,每个阶段一个 goroutine,channel 作为阶段之间的数据通道。

例子:数字生成 → 平方计算 → 打印结果,多个 channel 串联。

package feature

import (
	"fmt"
	"time"
)

// 生成数字,输出到channel
func gen(nums ...int) <-chan int {
	fmt.Println("[gen] 开始执行")
	out := make(chan int)
	go func() {
		fmt.Println("[G1 gen协程] 启动")
		for _, n := range nums {
			fmt.Printf("[G1] 准备发送 %d\n", n)
			out <- n
			fmt.Printf("[G1] 发送 %d 成功,休息2秒\n", n)
			time.Sleep(2 * time.Second)
		}
		fmt.Println("[G1] 准备关闭通道")
		close(out)
		fmt.Println("[G1] gen协程退出")
	}()
	fmt.Println("[gen] 返回通道c")
	return out
}

// square:延迟2秒再开始接收数据
func square(in <-chan int) <-chan int {
	fmt.Println("[square] 开始执行")
	out := make(chan int)
	go func() {
		fmt.Println("[G2 square协程] 启动")
		fmt.Println("[G2] 休眠结束,开始读取c通道")

		for n := range in {
			fmt.Printf("[G2] 收到 %d,计算平方\n", n)
			out <- n * n
			fmt.Printf("[G2] 发送平方结果 %d\n", n*n)
		}
		fmt.Println("[G2] 准备关闭out通道")
		close(out)
		fmt.Println("[G2] square协程退出")
	}()
	fmt.Println("[square] 返回通道out")
	return out
}

func ExecChannel() {
	fmt.Println("[main] start")
	c := gen(1, 2, 3, 4)
	fmt.Println("[main] 拿到c通道,准备调用square")
	out := square(c)
	fmt.Println("[main] 拿到out通道,开始range读取结果")
	for val := range out {
		fmt.Printf("[main] 收到结果:%d\n", val)
	}
	fmt.Println("[main] 程序结束")
}

// 输出结果
// Building .[main] start
// [gen] 开始执行
// [gen] 返回通道c
// [main] 拿到c通道,准备调用square
// [square] 开始执行
// [square] 返回通道out
// [main] 拿到out通道,开始range读取结果
// [G1 gen协程] 启动
// [G1] 准备发送 1
// [G2 square协程] 启动
// [G2] 休眠结束,开始读取c通道
// [G2] 收到 1,计算平方
// [G2] 发送平方结果 1
// [main] 收到结果:1
// [G1] 发送 1 成功,休息2秒
// [G1] 准备发送 2
// [G1] 发送 2 成功,休息2秒
// [G2] 收到 2,计算平方
// [G2] 发送平方结果 4
// [main] 收到结果:4
// [G1] 准备发送 3
// [G1] 发送 3 成功,休息2秒
// [G2] 收到 3,计算平方
// [G2] 发送平方结果 9
// [main] 收到结果:9
// [G1] 准备发送 4
// [G1] 发送 4 成功,休息2秒
// [G2] 收到 4,计算平方
// [G2] 发送平方结果 16
// [main] 收到结果:16
// [G1] 准备关闭通道
// [G1] gen协程退出
// [G2] 准备关闭out通道
// [G2] square协程退出
// [main] 程序结束

流程:gen chansquare chan,多个 channel 组合,形成流水线。

通道是无缓冲make(chan int),必须有另一个 goroutine 在接收,发送才会完成。