Go 多阶段 pipeline 必须用无缓冲或有界缓冲 channel 实现背压,避免死锁;各 stage 需独立 goroutine、显式 close 输入 channel、单向读写;错误传递需结合 error channel 或 context.Context 支持提前终止。

Go 里用 chan 实现多阶段 pipeline 的核心约束
必须用无缓冲或有界缓冲 chan,否则某个阶段卡住会直接拖垮整个流水线。无缓冲 channel 要求每个阶段都“同步就绪”,适合简单、低延迟场景;有界 channel(如 make(chan int, 100))能抗短时抖动,但容量设太大等于放弃背压控制,设太小又容易频繁阻塞。
常见错误现象:fatal error: all goroutines are asleep - deadlock——通常是因为某个 stage 没消费完上游发来的数据,或者忘了 close() 导致下游 range 永远等不到关闭信号。
- 每个 stage 都该是独立 goroutine,用
go func() { ... }()启动 - 输入 channel 必须被显式
close()(一般由源头或前一 stage 关闭),下游用for v := range inChan安全读取 - 不要在 pipeline 中间 stage 对同一 channel 既读又写——这是并发混乱的起点
怎么让 pipeline 支持错误传递和提前终止
原生 channel 不带错误语义,必须额外加一层控制通道。推荐用 error 类型的 channel 或复用 context.Context:前者更直观,后者利于超时和取消传播。
使用场景:某阶段解析 JSON 失败,后续阶段无需再处理;或用户主动取消请求,所有 stage 应尽快退出。
立即学习“go语言免费学习笔记(深入)”;
- 每个 stage 接收
ctx context.Context参数,并在 for 循环中检查select { case - 如果要用 error channel,约定只由第一个出错 stage 发送,其他 stage 监听该 channel 并立即 return(别再往下游发数据)
- 避免在 stage 内部 recover panic 后继续发送数据——这会让下游拿到脏值
pipeline 中如何安全地 fan-out / fan-in
fan-out 是把一个 channel 拆给多个并行 worker;fan-in 是把多个 channel 合并成一个输出。Go 标准库没提供现成工具函数,得手写,且必须注意 goroutine 泄漏和 channel 关闭时机。
参数差异:fan-out 数量影响吞吐和资源占用;fan-in 的 channel 数量决定 select 分支数,太多会拖慢调度。
- fan-out:用
for i := 0; i ,worker 内部用 <code>for v := range in - fan-in:用
func merge(cs ...,内部启动 goroutine 从每个 <code>c读并转发到返回的 channel;所有输入 channel 关闭后,才关闭输出 channel - 关键点:fan-in 的合并 goroutine 必须等待所有输入 channel 关闭,不能靠
len(cs)简单计数——因为 channel 可能提前 close
为什么 sync.WaitGroup 在 pipeline 里多数时候是错的
WaitGroup 适合“等所有 goroutine 结束”,但 pipeline 是流式处理,你并不知道何时“全部结束”——数据可能源源不断来,也可能中途断掉。强行 WaitGroup 会导致主流程死等,失去响应性。
性能影响:WaitGroup.Add/Wait 是原子操作,在高并发 stage 下有轻微开销;更严重的是逻辑错位——你真正该等的是“最后一个有效数据被消费完”,不是“最后一个 goroutine 退出”。
- 替代方案:用
context.WithCancel+ 显式关闭输入 channel,让每个 stage 自行退出 - 若真需确认“处理完毕”,应在最末端 stage 发送完成信号(如往
done chan struct{}写入),主 goroutine select 等待它 - WaitGroup 唯一合理用途:启动阶段的初始化 goroutine(比如加载配置),而非数据处理流水线本身
复杂点在于,pipeline 的关闭边界常常模糊——上游停了,但中间 stage 的缓冲区还有数据;下游关了,但上游还不知道。这些状态必须显式建模,不能依赖 goroutine 自然消亡。


















