Go 中流式计算依赖 chan 与 context.Context 组合,goroutine 仅作运行载体;错误做法是忽略 select/cancel 处理、不关闭 channel、复用 channel 导致阻塞或竞争。

Go 里没有“goroutine 实现流式计算”这种说法——goroutine 是并发单元,不是数据流抽象。真正支撑流式计算的是 chan 与 context.Context 的组合用法,goroutine 只是让每个 stage 能独立运行的载体。
为什么直接起 goroutine 不等于流式计算
常见错误是写这样的代码:go stage1(in) → go stage2(out1) → go stage3(out2),但 stage 函数内部仍是同步读写 channel,没加 select 和 ctx.Done(),结果一卡全卡。
- goroutine 启动后若不主动监听取消信号,就无法响应超时或中断
- channel 写操作不包在
select里,会永久阻塞 sender,导致上游 goroutine 挂住 - 忘记
defer close(out),下游range永远等不到 EOF,协程泄漏 - 复用同一个
chan作为输入和输出,造成读写竞争或逻辑混乱
每个 stage 必须是独立 goroutine + 只读/只写 channel
正确签名长这样:func filterErrors(ctx context.Context, in 。输入是只读、输出是只写,强制接口清晰。
- 函数内部必须用
go func() { ... }()启动,不能同步执行 - 输出 channel 要在 goroutine 结束前
defer close(out) - 所有
in读取和out写入都必须走select,包含分支 - 不要在 stage 内部
close(in)—— 输入由上游控制,你只消费
缓冲区大小不是性能开关,而是背压调节阀
用 make(chan T, 0)(无缓冲)能暴露阻塞点,适合调试;上线后按处理能力设缓冲,不是越大越好。
立即学习“go语言免费学习笔记(深入)”;
- IO 密集型 stage(如 HTTP 请求)建议缓冲 16–128,避免生产者狂推导致内存堆积
- CPU 密集型(如 JSON 解析、哈希计算)可设 512–2048,减少 goroutine 切换开销
- 绝对不要用
make(chan T, math.MaxInt)—— 这不是提速,是放弃对背压的感知 - 缓冲只是延缓阻塞,不能替代
ctx控制:即使有缓冲,也要在select里检查ctx.Done()
错误和取消信号必须走独立路径
别把错误塞进主数据 channel,否则消费者要不停判断 ok 和类型断言,极易出错。
- 推荐方案:每个 stage 返回两个 channel ——
和 <code>,由下游统一收集 - 更稳妥做法:返回
struct{ data T; err error },避免 channel 数量爆炸 - 用
errgroup.Group启动所有 stage,调用eg.Wait()可自动传播第一个 panic 或 cancel - worker pool 场景下,每个 worker 内必须
defer func(){ if r := recover(); r != nil { /* log & report */ } }()
最易被忽略的一点:很多人以为“起了 goroutine 就算流式”,其实关键在 channel 的生命周期管理与 context 的深度集成。一个没关的 out、一次漏掉的 select、一段没设缓冲的 IO 阻塞,都会让整条 pipeline 在线上悄无声息地卡死数小时。上线前用 go tool pprof -goroutine 看一眼阻塞数量,比补十次逻辑都管用。


















