Go Academy · 课时

流水线与阶段模式

组合多阶段并发流水线

第 1 / 4 课13 个步骤

流水线与阶段模式 是 CoddyKit 上的免费 Go Academy 课时。 这是第 1 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Go Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Go Academy 课程共包含 4 节课。

什么是流水线

并发流水线是一系列由通道连接的阶段。每个阶段从上游接收值,对其进行转换,然后将结果发送到下游,从而允许各阶段并发运行。

简单的三阶段流水线

生成 → 转换 → 消费:

func generate(nums ...int) <-chan int {
    out := make(chan int)
    go func() { defer close(out); for _, n := range nums { out <- n } }()
    return out
}
func square(in <-chan int) <-chan int {
    out := make(chan int)
    go func() { defer close(out); for n := range in { out <- n*n } }()
    return out
}

for v := range square(generate(2, 3, 4)) {
    fmt.Println(v) // 4, 9, 16
}

Stage 模式

每个 stage 都是一个接收输入通道并返回输出通道的函数。各 stage 可以自然组合,也可以独立并行化。

流水线中的扇出

将一个 stage 拆分为 N 个 goroutine,以并行执行受 CPU 限制的工作:

func parallelSquare(in <-chan int, n int) <-chan int {
    outs := make([]<-chan int, n)
    for i := 0; i < n; i++ {
        outs[i] = square(in) // all read from same in
    }
    return merge(outs...)
}

使用上下文取消

向每个 stage 传递一个上下文。上下文被取消时,各 stage 会停止读取并关闭其输出通道,从而排空整个流水线。

func stage(ctx context.Context, in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for {
            select {
            case v, ok := <-in:
                if !ok { return }
                select {
                case out <- process(v):
                case <-ctx.Done(): return
                }
            case <-ctx.Done(): return
            }
        }
    }()
    return out
}

背压

在各 stage 之间使用带缓冲的通道,以解除生产者和消费者速度之间的耦合。没有缓冲时,速度较慢的 stage 会阻塞整个流水线。

out := make(chan Result, 50) // buffer 50 results

错误传播

将值封装在结果结构体中,以便在流水线中传播错误而不会触发 panic:

type item struct { val int; err error }

Done 通道模式

对于较简单的流水线,可以使用 done 通道代替上下文。关闭 done 即可通知所有 stage 停止。

done := make(chan struct{})
defer close(done)

顺序保证

如果每个 stage 按顺序执行,简单流水线会保持顺序。并行化的 stage 不会保持顺序——如果需要,请使用索引字段在下游重新排列结果。

使用 WaitGroup 的流水线

最后一个 stage 不需要关闭通道,只需遍历该通道即可。在扇入操作中使用 WaitGroup,在所有生产者 goroutine 完成后关闭合并通道。

实际应用

流水线的常见应用包括:ETL 数据处理、图像缩放、日志解析、构建系统和流式请求处理程序。

快速检查

如何在多阶段流水线中传播取消信号?

回顾:流水线模式

要点:

  • Stage:func(in <-chan T) <-chan U——组合这些 stage 以构建流水线
  • 通过上下文实现取消;通过带缓冲的通道实现背压
  • 在 stage 内扇出以实现并行
  • 将值封装在结果结构体中以传播错误
免费开始

用 AI 导师学习 Go — 免费

在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。

课程
51
课程
203

常见问题解答

「流水线与阶段模式」课时是免费的吗?

是的 — 「流水线与阶段模式」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Go Academy 课程的其余内容,请升级到 CoddyKit PRO。 Go Academy 课程共包含 4 节课。

「流水线与阶段模式」这节课中我会学到什么?

组合多阶段并发流水线 你通过在浏览器中直接运行的动手代码来练习 Go Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Go Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Go Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。

「流水线与阶段模式」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Go Academy 课中编写并运行代码吗?

能。每节 Go Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 流水线与阶段模式
  2. 使用 errgroup 处理并发错误
  3. 信号量模式
  4. 泄漏检测与取消
← 返回 Go Academy