go多阶段pipeline错误传播需channel关闭与error返回双机制,配合context.withcancel统一控制生命周期,各stage监听ctx.done()并及时关闭输入输出channel,避免goroutine泄漏和数据丢失。

Go 中多阶段 pipeline 的错误传播必须靠 channel 关闭 + error 返回双机制
单靠 chan error 或只在最后汇总错误,都会导致中间阶段 panic、goroutine 泄漏或数据丢失。真实场景里,一个 stage 失败后,后续 stage 必须停止消费、尽快退出,且上游不能继续往 pipeline 塞数据。
典型错误是只用 select 监听 errCh,却没同步关闭输入 channel —— 导致下游 goroutine 卡在 range 里永远等不到 close。
- 每个 stage 启动时,都应接收一个
ctx.Context,并在收到ctx.Done()时立即退出 - 所有输入 channel 都要带缓冲(如
make(chan Item, 16)),避免上游因下游阻塞而卡死 - 错误发生后,先调用
cancel(),再关闭所有输出 channel;不能反过来,否则下游可能读到部分数据后突然被 close 中断
用 context.WithCancel 控制 pipeline 全局生命周期
不用 context.WithCancel,就无法让所有 stage 对同一个信号响应。手动传 done chan struct{} 容易漏传、重复 close,且无法与超时、取消链路集成。
关键点在于:cancel 函数必须由第一个失败 stage 调用,且只能调一次;其他 stage 在 select 中监听 ctx.Done() 并 return,不自行 cancel。
- 主函数创建
ctx, cancel := context.WithCancel(context.Background()),把ctx传给每个 stage - 每个 stage 的 goroutine 开头写
defer func() { recover() }()是错的 —— panic 不等于错误,不应掩盖业务错误 - stage 内部若调用可能阻塞的外部服务(如 HTTP 请求),必须把
ctx传进去,否则 timeout 不生效
每个 stage 必须显式处理“上游已关闭”和“ctx 已取消”两种退出条件
只检查 ok := 不够 —— 这只能捕获上游 close,但若 ctx 先 cancel,channel 还开着,goroutine 就会 hang 住。
正确模式是始终用 select 双路监听,且 default 分支不可省略(防止 busy loop):
for {
select {
case item, ok :=
- 不要在
case item, ok := 里直接做耗时操作 —— 若此时 ctx 已 cancel,会浪费资源 - 若 stage 需要批量处理(如 collect 10 个再 send),必须在每次
select前检查ctx.Err() != nil,否则可能攒满 buffer 后才响应 cancel - 输出 channel 的 send 操作也必须在
select里做,否则可能因下游阻塞而卡死
error 信息需附带 stage 名和原始 error,不能只返回 fmt.Errorf("failed")
线上排查时,只知道“pipeline 失败”毫无价值。必须明确是哪个 stage、在处理哪个数据时、因为什么底层 error(如 io.EOF、json.SyntaxError)失败。
建议统一用结构体包装错误,而不是拼字符串:
type PipelineError struct {
Stage string
ItemID string
Err error
}
func (e *PipelineError) Error() string {
return fmt.Sprintf("stage %s item %s: %v", e.Stage, e.ItemID, e.Err)
}
- 每个 stage 在 return 前,应 new 一个
*PipelineError,把当前 stage 名、当前 item 的标识(如 ID 或 hash)、原始 error 全部带上 - 不要用
errors.Wrap层层套 —— 日志里会看到冗长的 stack,但真正需要的是 stage 上下文,不是调用栈 - 如果某个 stage 本身不产生新 error(如只是 map 转换),但上游传来的 item 已含 error 字段,应原样透传,不要吞掉
最麻烦的其实是测试 —— 很多人只测 happy path,结果上线后某个 stage 因网络抖动提前退出,整个 pipeline 就静默卡死。一定要写强制某个 stage 返回 error 的单元测试,验证 cancel 是否广播、goroutine 是否回收、channel 是否全部关闭。这些细节不验证,等于没做错误传播。











