goroutine流水线的核心是“分阶段+通道串联”,即把数据处理拆分为多个goroutine阶段,各阶段通过只读/只写channel串联,关键在于精确控制通道关闭时机与阻塞行为,避免死锁或panic。

goroutine流水线的核心是“分阶段+通道串联”
Go语言里没有现成的“流水线”语法糖,所谓流水线就是把数据处理拆成多个 goroutine 阶段,每个阶段通过 chan 接收输入、处理、再发给下一阶段。关键不是开多少 goroutine,而是控制好通道的关闭时机和阻塞行为——否则极易卡死或 panic。
典型错误现象:fatal error: all goroutines are asleep - deadlock,基本都源于某个阶段没关 channel,或者某处用 range 读取未关闭的 channel。
- 每个阶段应独立启动 goroutine,主 goroutine 不参与处理逻辑
- 输入 channel 建议设为只读(
<-chan T),输出 channel 设为只写(chan<- T),避免误写/误读 - 若某阶段需扇出(fan-out)多个 worker,必须用
sync.WaitGroup等待全部完成后再关闭下游 channel
如何正确关闭流水线中的 channel
channel 关闭必须由“生产者”负责,且只能关一次;消费者不能关,也不能假设上游一定关闭。常见反模式是让中间 stage 主动关 output channel,结果下游还在 range 就 panic。
正确做法:上游 stage 在所有数据发送完毕后关闭自己的 output channel;下游 stage 在收到关闭信号后退出循环,并可选择关闭自己的 output channel(如果它是最后一级或确定不再产出)。
立即学习“go语言免费学习笔记(深入)”;
- 用
for v := range in自动处理关闭,比for { select { case v, ok := <-in: if !ok { break } }}更简洁安全 - 若 stage 有多个输入 channel(如合并两个源),需用
select+ok判断各 channel 是否关闭,不能依赖range - 不要在 defer 中关闭 output channel——可能在 goroutine 启动前就执行了
扇入(fan-in)和扇出(fan-out)怎么写才不漏数据
扇出适合 CPU 密集型任务并行化(如解析 JSON、校验字段),扇入适合聚合多个来源(如合并日志流、合并 API 响应)。二者都依赖 channel 复用,但容易因 goroutine 泄漏或 channel 缓冲不足丢数据。
示例:3 个 worker 并行处理数据,结果统一汇入一个 channel:
func fanOut(in <-chan int, workers int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for v := range in {
out <- v * v // 模拟处理
}
}()
}
go func() {
wg.Wait()
close(out)
}()
return out
}- 务必用
sync.WaitGroup等待所有 worker 结束,再关out,否则主 goroutine 可能提前退出 - 输入 channel
in必须由调用方保证关闭,否则 worker 会永远阻塞在range - 避免给
out设缓冲(如make(chan int, 100)),除非明确需要背压控制;否则易掩盖阻塞问题
为什么加 buffer 容易掩盖设计缺陷
有人一遇到 “goroutine 卡住” 就给 channel 加 buffer,比如 make(chan int, 1000)。这看似解决问题,实则把阻塞延迟暴露、掩盖了真正的吞吐瓶颈或逻辑错误。
buffer 的本质是临时队列,它不解决“谁来消费”“谁来关闭”“下游是否跟得上”的问题。一旦下游处理变慢,buffer 会积压,内存上涨,最终 OOM 或超时。
- 仅在明确需要解耦发送/接收节奏时用 buffer,例如日志采集中暂存几秒数据
- buffer 大小 ≠ 并发数,更不等于吞吐量——真实吞吐取决于最慢 stage 的处理速度
- 调试阶段建议一律用无缓冲 channel(
make(chan int)),能最快暴露阻塞点
流水线的复杂性不在 goroutine 数量,而在 channel 生命周期的协作关系。一个 stage 忘记关 channel,整条链就僵住;一个 stage 过早关 channel,下游就收不到后续数据。这些边界条件,没法靠加 goroutine 解决。


















