Go pipeline 模式
用 channel 串联多阶段处理,明确数据流、fan-out/fan-in、关闭顺序与取消传播;对齐 go.dev/blog/pipelines。
[!info] 关联笔记
Go pipeline 模式
这个概念为什么会出现
许多工作不是“一堆彼此无关的任务”,而是同一批数据依次经过多个阶段:
生成 → 变换 → 过滤 → 聚合 → 输出
若把所有阶段塞进一个函数,难测、难并行、难取消。若每阶段一个 goroutine、用 channel 连接,则:
- 阶段边界变成类型边界
- 某阶段可 fan-out 并行
- 取消与关闭可沿管道传播
这正是 Go Concurrency Patterns: Pipelines 的主题:pipeline 是 channel 上的阶段组合,而不是某个框架类名。
[!abstract] 一句话理解 pipeline 把处理拆成“接收输入 channel → 发送输出 channel”的阶段;数据向前流,关闭表示流结束,context 表示该停,fan-out/fan-in 调节并行与汇合。
最小可运行示例
先把示例放进业务场景,再看代码:
场景:指标流水线——采集 → 平方放大 → 再平方(演示级联)
监控 agent 里常见:
- gen:从配置列出原始样本
2, 3 - sq:每一级做变换(这里用平方假装“归一化/放大”)
- 下游
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
结合场景再看三个关注点
- 每阶段返回
<-chan,隐藏发送与关闭,所有权清晰。 - 关闭沿下游传播:上游 close → 本段 range 完 → close 自己的 out。
- 真实系统还要 ctx 取消与错误 channel;本最小例只钉住数据流骨架。
核心概念与准确模型
阶段函数的形状
官方博文推荐的风格:
func stage(in <-chan T) <-chan U
// 或
func stage(ctx context.Context, in <-chan T) <-chan U
阶段内部:
- 创建
out - 启动 goroutine:读
in,写out 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 的分工
| pipeline | worker pool | |
|---|---|---|
| 结构 | 多阶段数据流 | 单队列多工人 |
| 重点 | 阶段边界与组合 | 并发度上限 |
| 自然形态 | 类型化 channel 链 | jobs/results + N workers |
| 可嵌套 | 某阶段内部可用 pool | pool 任务内可再跑子 pipeline |
二者常组合:pipeline 的重阶段内部 fan-out 成 pool。
取消与泄漏(博文后半的重点)
若下游提前返回(只要第一个结果),上游仍可能阻塞在发送上 → goroutine 泄漏。正确做法:
- 全程传
ctx - 发送用
select:case out <- v:/case <-ctx.Done(): return - 确保所有阶段在取消时退出并 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 原生传的是值流。错误策略:
- 值内嵌
struct{ V T; Err error } - 错误即取消:
errgroup取消 ctx,阶段听 Done - 副作用日志:仅可观测,不替代返回
避免“错误 channel + 数据 channel”双通道却无统一关闭协议。
边界情况
- 下游不消费:无缓冲链路上游全堵;无取消则泄漏。
- 某阶段 panic:整链断裂;需 recover 策略或让进程崩溃可观测。
- 缓冲过大:掩盖背压,内存隐性增长。
- 阶段过多:调度与 channel 开销超过并行收益。
- 共享可变消息:channel 传指针时所有权仍要约定。
- 把分布式消息队列当本地 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。
工程实践
- 阶段纯净:输入输出 channel + ctx;副作用集中在边缘阶段。
- 类型化消息:比
chan any更易推理。 - 可测性:单阶段单测(给定 in 序列,断言 out 序列)。
- 观测:每阶段计数、滞留时间、丢弃/错误数。
- 背压优先于丢弃;若必须丢弃,指标与策略显式化。
- 与 HTTP/gRPC 结合:handler 的 ctx 注入首阶段;客户端断开即 cancel 整链。
- 先读官方 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 与关系篇串生命周期。
自测题
概念题
- pipeline 阶段函数为什么常返回
<-chan T? - fan-out 与 fan-in 各解决什么问题?
- 为什么 pipelines 博文强调 cancellation?
代码推理题
下游在 for range out 中 break 后,无缓冲上游卡在 out <- v 且无 ctx。会发生什么?
工程思考题
日志采集:读文件 → 解析 JSON → 过滤 → 批量写远端。哪些阶段值得并行?取消从哪里注入?
参考答案
展开
- 隐藏发送与 close 所有权,组合更安全。
- fan-out 提高可并行阶段吞吐;fan-in 安全汇合多路输出。
- 下游可能不再消费;无取消则发送阻塞,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 pool 或 Go context
- 本篇状态:已深化(对齐官方 pipelines 博文)