结论:用 func(input) 可直接调用函数处理输入;需确保 input 类型与 func 参数要求一致,否则可能报错或产生意外结果。

直接说结论:用 func(in 形式定义中间件函数,每个函数启动独立 goroutine 处理数据流,靠 channel 串接,不共享状态、不阻塞上游。
中间件函数必须返回 并启动 goroutine
常见错误是写成同步函数,比如 func(in []int) []int 或直接在调用方 goroutine 里循环处理——这会把并发逻辑“压平”,失去流水线意义。
- 输入参数必须是只读 channel:
in ,表示“只从这里读” - 返回值必须是只写 channel:
out chan,但更推荐返回 <code>(只读),避免下游误写 - 函数体里必须用
go func() { ... }()启动新 goroutine,否则调用即阻塞 - 别在函数内直接
close(in)—— 输入 channel 只能由上游关,中间件只负责关自己的out
如何安全关闭输出 channel
没关 out 是 pipeline 卡死的最常见原因:下游 for range out 永远等不到 close 信号。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
- 关闭时机:输入 channel 关闭后,处理完所有剩余数据,再
close(out) - 必须在 goroutine 内部关,不能 defer 在外层函数——外层 return 时 goroutine 可能还没跑完
- 错误写法:
defer close(out)放在函数开头;正确写法是defer close(out)放在 goroutine 的末尾 - 如果阶段可能提前退出(如过滤器跳过某些项),仍要确保
close(out)执行,哪怕没发任何数据
支持错误传播时,别用裸 channel
原生 chan int 无法区分“流结束”和“出错了”。一旦某阶段 panic 或返回 error,下游只能卡住或 panic。
- 推荐统一用带错误的结构体:
type Result struct { Val int; Err error } - 所有中间件函数签名改为:
func(in - 任意阶段发现错误,立即
out ,下游收到后可选择终止整个流 - 不要依赖
recover()捕获 panic 来传递错误——它跨不了 goroutine,且掩盖真实崩溃点 - 更健壮的做法是结合
context.Context:在每个 goroutine 中 selectctx.Done(),提前退出并 close out
缓冲 channel 要谨慎设大小
无缓冲 channel(make(chan int))会让上下游严格同步,容易因下游慢导致上游卡住;缓冲 channel(make(chan int, 10))能缓解,但不是万能解药。
- 缓冲区大小 ≠ 并发度,只是临时队列容量。设太大可能吃光内存,太小起不到解耦作用
- 扇出场景(如 1 个输入 → 3 个 worker)建议对 worker 输入 channel 设缓冲,比如
make(chan int, 100) - 扇入场景(多个 worker → 1 个输出)必须用无缓冲或小缓冲,否则可能丢数据或乱序
- 生产环境建议配合背压策略:当
len(ch) == cap(ch)时主动等待或丢弃,而不是盲目写入
真正难的不是写几个 go func(),而是谁关 channel、什么时候关、关错会怎样、错误怎么透出——这些细节不抠清楚,pipeline 看似跑通,实则随时可能在线上静默卡死或 panic 泄露 goroutine。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










