Go pipeline 模式

用 channel 串联多阶段处理,明确数据流、fan-out/fan-in、关闭顺序与取消传播;对齐 go.dev/blog/pipelines。

#type / concept #status / growing #tech / dev #resource / go

[!info] 关联笔记

Go pipeline 模式

这个概念为什么会出现

许多工作不是“一堆彼此无关的任务”,而是同一批数据依次经过多个阶段

生成 → 变换 → 过滤 → 聚合 → 输出

若把所有阶段塞进一个函数,难测、难并行、难取消。若每阶段一个 goroutine、用 channel 连接,则:

  • 阶段边界变成类型边界
  • 某阶段可 fan-out 并行
  • 取消与关闭可沿管道传播

这正是 Go Concurrency Patterns: Pipelines 的主题:pipeline 是 channel 上的阶段组合,而不是某个框架类名。

[!abstract] 一句话理解 pipeline 把处理拆成“接收输入 channel → 发送输出 channel”的阶段;数据向前流,关闭表示流结束,context 表示该停,fan-out/fan-in 调节并行与汇合。

最小可运行示例

先把示例放进业务场景,再看代码:

场景:指标流水线——采集 → 平方放大 → 再平方(演示级联)

监控 agent 里常见:

  1. gen:从配置列出原始样本 2, 3
  2. sq:每一级做变换(这里用平方假装“归一化/放大”)
  3. 下游 for range 消费最终流

pipeline 的要点不是“平方”,而是阶段形状

  • 入:<-chan T
  • 出:<-chan U
  • 内部 goroutine 发送并在结束时 close(out)
  • 调用方只组合阶段,不碰关闭细节
package main

import "fmt"

// emitSamples 模拟“采集阶段”:产出原始指标。
// 返回只读 channel:调用方不能误 close/误发送。
func emitSamples(nums ...int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out) // 发送方关闭 = 没有更多样本
		for _, n := range nums {
			out <- n
		}
	}()
	return out
}

// square 模拟“变换阶段”:读上游,写出变换后的值。
// 上游关闭 → range 结束 → 关闭自己的 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
}

func main() {
	// 组合:样本 → 平方 → 再平方;主流程只消费最终流
	// 2 → 4 → 16;3 → 9 → 81
	for n := range square(square(emitSamples(2, 3))) {
		fmt.Println("metric:", n)
	}
}

建议运行:

go run .

期望输出:

metric: 16
metric: 81

结合场景再看三个关注点

  1. 每阶段返回 <-chan,隐藏发送与关闭,所有权清晰。
  2. 关闭沿下游传播:上游 close → 本段 range 完 → close 自己的 out。
  3. 真实系统还要 ctx 取消与错误 channel;本最小例只钉住数据流骨架。

核心概念与准确模型

阶段函数的形状

官方博文推荐的风格:

func stage(in <-chan T) <-chan U
// 或
func stage(ctx context.Context, in <-chan T) <-chan U

阶段内部:

  1. 创建 out
  2. 启动 goroutine:读 in,写 out
  3. in 耗尽后 close(out)

调用方只组合 channel,不碰关闭细节。

关闭传播

gen close(out1) → stageA range 结束 close(out2) → stageB ...

规则复述:

  • 只由发送方 close
  • 多发送方写入同一 channel 时,必须 fan-in 或 WaitGroup 后再 close
  • close 表示“无更多值”,不是取消进行中计算的充分条件

fan-out / fan-in

来自 pipelines 博文的核心术语:

模式含义典型用途
fan-out多个 stage 实例读同一输入并行重计算阶段
fan-in多个输入合并到一个 channel汇总并行结果
// 伪代码结构
in := gen(...)
c1 := sq(in)
c2 := sq(in) // fan-out:两个 sq 竞争读 in(注意 in 只能被安全地多读若语义允许)
for n := range merge(c1, c2) { // fan-in
    ...
}

merge 必须在所有输入关闭后才 close 输出。

与 worker pool 的分工

pipelineworker pool
结构多阶段数据流单队列多工人
重点阶段边界与组合并发度上限
自然形态类型化 channel 链jobs/results + N workers
可嵌套某阶段内部可用 poolpool 任务内可再跑子 pipeline

二者常组合:pipeline 的重阶段内部 fan-out 成 pool。

取消与泄漏(博文后半的重点)

若下游提前返回(只要第一个结果),上游仍可能阻塞在发送上 → goroutine 泄漏。正确做法:

  1. 全程传 ctx
  2. 发送用 selectcase out <- v: / case <-ctx.Done(): return
  3. 确保所有阶段在取消时退出并 close 自己拥有的 out(或文档化由谁 close)
func sq(ctx context.Context, in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for n := range in {
			select {
			case out <- n * n:
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

错误如何流动

pipeline 原生传的是值流。错误策略:

  1. 值内嵌 struct{ V T; Err error }
  2. 错误即取消errgroup 取消 ctx,阶段听 Done
  3. 副作用日志:仅可观测,不替代返回

避免“错误 channel + 数据 channel”双通道却无统一关闭协议。

边界情况

  1. 下游不消费:无缓冲链路上游全堵;无取消则泄漏。
  2. 某阶段 panic:整链断裂;需 recover 策略或让进程崩溃可观测。
  3. 缓冲过大:掩盖背压,内存隐性增长。
  4. 阶段过多:调度与 channel 开销超过并行收益。
  5. 共享可变消息:channel 传指针时所有权仍要约定。
  6. 把分布式消息队列当本地 pipeline:进程内 channel 语义不同(无持久化、无消费组)。

常见误区

[!warning] 常见误区:pipeline 一定更快 错误:三阶段拆分必然提速。
正确:拆分首先为清晰与可取消;加速来自可并行阶段与匹配的 fan-out。

[!warning] 常见误区:关闭 channel 等于取消 错误:close 后以为所有计算停止。
正确:close 结束流;正在跑的工作要靠 ctx 或可中断调用。

[!warning] 常见误区:每个阶段无界缓冲 错误:make(chan T, 1e9) 防阻塞。
正确:有界 + 背压 + 取消;内存是一等资源。

[!warning] 常见误区:多发送者随意 close 错误:fan-out 后谁先结束谁 close。
正确:merge/协调者关闭;见 go-channel-ownership-and-closing

工程实践

  1. 阶段纯净:输入输出 channel + ctx;副作用集中在边缘阶段。
  2. 类型化消息:比 chan any 更易推理。
  3. 可测性:单阶段单测(给定 in 序列,断言 out 序列)。
  4. 观测:每阶段计数、滞留时间、丢弃/错误数。
  5. 背压优先于丢弃;若必须丢弃,指标与策略显式化。
  6. 与 HTTP/gRPC 结合:handler 的 ctx 注入首阶段;客户端断开即 cancel 整链。
  7. 先读官方 pipelines 博文再发明框架——多数需求是组合函数,不是新运行时。

可验证实验

实验 1:关闭链

gen → sq → print,确认只需 gen 关闭,后续自动收尾。

实验 2:提前退出泄漏

构造下游只读一个值就 return、上游大发送;对比有无 ctx 的 goroutine 数量(runtime.NumGoroutine)。

实验 3:fan-in

两路 sq 写入 merge,检查输出个数与是否死锁。

实验 4:与 pool 对比

同一批“无阶段关系”任务分别用 pool 与硬套三阶段 pipeline,比较代码复杂度。

本节总结

  • 本质:用 channel 连接的阶段式数据流。
  • 关键规则:发送方关闭、取消防泄漏、fan-out/fan-in 的关闭纪律。
  • 最易错:下游提前返回导致上游阻塞;关闭权争夺。
  • 下一步worker pool 控并发;context 与关系篇串生命周期。

自测题

概念题

  1. pipeline 阶段函数为什么常返回 <-chan T
  2. fan-out 与 fan-in 各解决什么问题?
  3. 为什么 pipelines 博文强调 cancellation?

代码推理题

下游在 for range outbreak 后,无缓冲上游卡在 out <- v 且无 ctx。会发生什么?

工程思考题

日志采集:读文件 → 解析 JSON → 过滤 → 批量写远端。哪些阶段值得并行?取消从哪里注入?

参考答案

展开
  1. 隐藏发送与 close 所有权,组合更安全。
  2. fan-out 提高可并行阶段吞吐;fan-in 安全汇合多路输出。
  3. 下游可能不再消费;无取消则发送阻塞,goroutine 泄漏。
    代码题:上游永久阻塞,对应 goroutine 泄漏;整链可能无法退出。
    工程题:解析/过滤可 fan-out;写远端受连接与配额限制宜有界;取消从进程信号或上游 ctx 注入首阶段并贯穿。

延伸阅读与资料来源

资料类型支撑
Go Concurrency Patterns: Pipelines官方博客模式本体、fan-out/in、取消
Go Concurrency Patterns: Context官方博客取消树
Advanced Go Concurrency Patterns演讲/博客更深并发套路
Share Memory By Communicating博客channel 设计示例
Spec — Channel types规范基础语义

笔记元信息

  • 建议文件名:go-pipeline-pattern.md
  • 所属阶段:并发模式
  • 建议下一篇:Go worker poolGo context
  • 本篇状态:已深化(对齐官方 pipelines 博文)
创建于 2026/6/20 更新于 2026/7/15