流水线卡在中间阶段主因是输出 channel 未关闭或消费者提前退出致上游阻塞;for range 可自动感知 close(in) 退出,而手动循环需用 ok 判断。

Go 流水线不是加几个 go 就能提速的——没处理好 channel 关闭、goroutine 生命周期和背压,反而会让程序卡死或内存暴涨。
为什么流水线卡在中间阶段不往下走?
最常见原因是某个 stage 的 goroutine 没关输出 channel,或者消费者提前退出,导致上游写 out 时被同步 channel 阻塞住。
- 用
for range in读输入 channel,它会自动感知close(in)并退出循环;手动for { select { case x, ok := 容易漏掉关闭信号 - 每个 stage 必须在处理完所有输入后调用
close(out),不能靠 defer(除非确保 goroutine 能跑完) - 避免在 stage 内部再起 goroutine 往
out写——关闭时机失控,下游永远等不到 EOF
如何让多个 worker 并行处理又不破坏顺序或并发写冲突?
扇出(fan-out)适合 CPU 密集型中间阶段,但要注意输出聚合方式。
- 需要保序:用单个 goroutine 从多个 worker 的 output channel 中
select读取,按需排序或直接转发(如日志流水线) - 需要聚合统计(如词频):worker 只发结果到一个
chan map[string]int,由单独的汇总 goroutine 合并 map,避免并发写 panic - 不要让多个 worker 直接往同一个
map写——即使加了sync.Mutex,也违背流水线“解耦”初衷,还容易成为瓶颈
context.Context 该传给谁、怎么用才真正生效?
Context 不是装饰,传错位置等于没传:cancel 无效、goroutine 泄露、超时判断失准。
立即学习“go语言免费学习笔记(深入)”;
- 必须作为参数传进每个 stage 函数,例如
func filter(ctx context.Context, in - stage 内部要用
select主动监听ctx.Done(),不能只依赖range in自然退出(上游可能还在发数据) - 若 stage 内发起 HTTP 请求,必须用
http.NewRequestWithContext(ctx, ...),否则 cancel 对它完全无效 - 别在 stage 里用
context.WithTimeout(ctx, time.Second)覆盖父 context——外部调用CancelFunc就失效了
缓冲区大小设多少才不算拍脑袋?
缓冲太小,生产者频繁阻塞;缓冲太大,内存占用不可控,还掩盖消费慢的问题。
- IO 密集型阶段(如读文件、HTTP 请求)建议用带缓冲 channel,容量 ≈ 预估并发请求数 × 单次响应平均大小(字节)/ 1024,向上取整到 64 或 128
- CPU 密集型阶段(如解密、压缩)缓冲可设小些(如 8–32),靠 worker 数量调节吞吐,避免堆积大量待处理数据
- 关键链路加限速:用
time.Ticker控制每秒最多进 N 条,比无脑扩缓冲更可控
真正难的从来不是拼出几个 chan 和 go,而是每个 stage 的输入契约是否清晰、关闭时机是否确定、错误是否可传播、context 是否穿透到底——这些细节没对齐,流水线就只是看起来像流水线。


















