Go 的 Pipeline 需显式用 goroutine 和 channel 编排,关闭时机错误、忽略背压或漏传 context.Context 均致卡死或泄露;for v := range in 卡住因上游未关 in 或中间 stage 未关 out;range 阻塞等待 channel 关闭,仅发送方可 close(out) 且不可 defer;过滤型 stage 须处理完所有输入后 close(out);禁用无缓冲 select 循环以防背压失控。

Go 的 Pipeline 不是语法糖,而是靠显式 goroutine + channel 编排出来的数据流;用错关闭时机、忽略背压、漏传 context.Context,三者任一都会让整条流水线卡死或泄露。
为什么 for v := range in 会卡住?
根本不是循环写错了,而是上游没关 in,或中间 stage 没关自己的 out。Go 的 range 会一直阻塞等待新值,直到 channel 被关闭 —— 如果没人关,它就永远等下去。
- 只允许发送方(即启动 goroutine 并往
out写的那一个)调用close(out),且只能调一次 - 别在 stage 函数里用
defer close(out):函数 return 时 goroutine 可能还没启动,defer就白挂了 - 如果 stage 做过滤(比如跳过奇数),仍要确保处理完所有输入后,
close(out)—— 否则下游range永不退出 - 避免用
for { select { case v, ok := :<code>break只跳出当前 select,不会退出 for,容易漏关
如何选缓冲大小:无缓冲 vs make(chan T, N)
缓冲不是“越大越好”,它是背压控制的开关。设错直接触发 fatal error: all goroutines are asleep - deadlock。
- 无缓冲 channel(
make(chan int)):收发必须同步就绪,适合低延迟、阶段处理耗时稳定的小流量场景 - 有界缓冲(如
make(chan int, 100)):抗短时抖动,但容量 > 实际峰值吞吐 × 处理延迟,等于放弃背压; - 绝对不要用
make(chan int, 0x7fffffff)或大 slice 模拟“无限缓冲”:内存爆掉前 pipeline 先假死 - 若不确定,先从
make(chan T, 16)起手,压测时观察 goroutine 数和 channel 阻塞率再调
怎么让 pipeline 支持错误中断和超时?
原生 channel 没错误语义,close() 只代表“数据结束”,不区分成功/失败。靠 panic 传播?跨 goroutine 无效。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
立即学习“go语言免费学习笔记(深入)”;
- 每个 stage 必须接收
ctx context.Context参数,并在主循环中用select监听ctx.Done() - 错误统一走单独的
errCh chan error,约定仅第一个出错 stage 发送,其他 stage 收到后立即return,不再读in、也不再写out - 别在 stage 里
recover()吞 panic:掩盖真实崩溃点,且无法通知下游终止 - 若某 stage 调外部服务(如 HTTP),必须把
ctx传进去,用http.NewRequestWithContext(ctx, ...),否则超时取消不生效
fan-in 和 fan-out 怎么不泄漏 goroutine?
扇入扇出不是加几个 go 就完事,合并 goroutine 必须等所有输入 channel 关闭才关输出,否则下游永远卡在 range。
-
fan-out:用for i := 0; i ,<code>worker内部必须for v := range in+defer close(out)(注意:这里是out,不是in) -
fan-in:不能靠len(cs)计数,要用sync.WaitGroup或为每个输入 channel 启一个 goroutine,全部结束后才close(merged) - 所有参与
fan-in的输入 channel 必须被显式关闭(由各自上游负责),否则合并 goroutine 会永远等不到 EOF - 避免在
fan-in合并逻辑里做耗时操作(如 JSON 解析):它本该是轻量转发层,重逻辑应下沉到 worker stage
最常被跳过的细节是:每个 stage 的 out 关闭时机,必须严格绑定于“输入消费完毕 + 所有输出已发送”,而不是函数返回、goroutine 结束或随便一个 defer。少这一步,整条 pipeline 就是纸糊的。

















