加buffer channel易出错是因为它掩盖背压缺失,导致内存暴涨或数据丢失;仅在上游速率稳定、下游能力可预测且延迟敏感时才适用,IO密集型用小缓冲(16~128),CPU密集型可用512~2048,严禁math.MaxInt,且每个stage必须显式close(out)否则下游range永久阻塞。

为什么 pipeline 中间加 buffer channel 容易出错
加 make(chan int, 10) 看似能缓解阻塞,但会掩盖背压缺失问题,导致内存暴涨或数据丢失。
使用场景:仅当明确知道上游生产速率稳定、下游消费能力可预测,且延迟敏感(如实时日志采样)时才考虑缓冲。
- IO 密集型阶段(如 HTTP 请求、DB 查询)建议小缓冲(
16~128),让生产者稍作等待,避免下游积压 - CPU 密集型阶段(如 JSON 解析、哈希计算)可适度加大(
512~2048),减少调度开销 - 永远不要用
make(chan T, math.MaxInt)—— 这不是优化,是放弃控制 - 缓冲 channel 不改变关闭语义:下游仍需靠
range检测关闭信号,否则无法退出循环
每个 stage 必须显式 close(out) 吗
必须。不关 out,下游 for x := range in 就永远卡住——Go 的 channel 是同步的,send 会等 recv 准备好;如果没人读、也没关,上游就永久阻塞在 out 。
- 关闭动作应在 goroutine 内部完成,且必须在所有数据发送完毕后:用
defer close(out)最安全 - 别在 stage 外部(比如 main 函数)关
out:违反职责分离,且容易关早(数据还没发完)或关晚(goroutine 已退出) - 生成器(source)关自己输出;处理器(transform)关自己输出;消费者(sink)不关任何 channel,只消费
- 若 stage 内部有扇入(多个输入 channel),需用
sync.WaitGroup或select+done保证所有输入读完再关out
如何让多个 stage 严格按顺序执行且不丢数据
靠 goroutine 调度保序不可行。真要顺序,就得放弃“多 stage 并发写同一 channel”,改用主协程显式控制消费节奏。
立即学习“go语言免费学习笔记(深入)”;
- 为每个 stage 分配独立 channel:
headCh、bodyCh、footCh - 主协程按需依次
for x := range headCh→for x := range bodyCh→for x := range footCh - 每个 stage 启动后立即关闭自己的
out,确保主协程不会因某 stage 卡住而阻塞后续阶段 - 若某 stage 出错,应向统一
errCh chan error发送错误,并提前关闭其out,主协程 select 到 err 后直接 break
context.Context 该传给谁、怎么传才有效
Context 不是装饰品,传错位置会导致 cancel 无效、goroutine 泄露或超时判断失准。
- 必须作为参数传入每个 stage 函数签名,例如:
func square(ctx context.Context, in - stage 内部所有阻塞操作(
http.Do、db.QueryRow、time.Sleep)都应使用该 ctx,而非硬编码新 context - 监听
ctx.Done()应放在循环内,且优先级高于 channel 接收:select { case - 一旦收到
ctx.Done(),stage 应立即停止处理、关闭out,并返回——不等待未完成的子 goroutine,由上层用sync.WaitGroup统一回收
真正卡住 pipeline 的,从来不是逻辑复杂度,而是 channel 关闭时机和 context 传播路径这两个细节。前者决定数据流是否完整,后者决定系统能否及时响应中断。写完一个 stage,先问自己:它关了 out 吗?它响应 ctx.Done() 吗?它有没有把错误传出去?这三个问题没答清楚,流水线就算跑通,也只是侥幸。



















