goroutine + channel 的乐高 — worker pool / pipeline / fan-out fan-in / errgroup 首错取消 / 信号量限流: 谁启动, 谁收尾
goroutine 和 channel 是乐高积木, 并发模式就是几种标准拼法: worker pool 用"N 个 worker 抢一个 jobs chan"把并发钉在上限上; pipeline 用"每阶段一个 goroutine + 一个 chan"让数据像流水线一样流过去, 内存里永远只有一小段; fan-out/fan-in 负责扇开再收回一口。真正让它们能在生产活下来的是两条纪律: 限流(信号量或 errgroup.SetLimit, 控并发就是控内存) 和 收口(errgroup 首错取消 + 谁启动谁收尾)。模式本身不难, 难的是每条 goroutine 的退出路径都想清楚。
jobs := make(chan Job, 200) results := make(chan Result) for i := 0; i < 8; i++ { go worker(jobs, results) } // 关键: 8 个 worker = 并发上限, 200 只是排队缓冲
close 驱动下游 for range 结束, 像流水线传桶, 任意时刻内存里只有在制品。func double(in <-chan int) <-chan int { out := make(chan int) go func() { defer close(out) // 谁开 chan 谁收尾 for v := range in { out <- v * 2 } }() return out // 上游 close → 本 range 结束 → close 传导 }
wg.Wait() 之后才 close(out) —— 必须等所有写端退出。var wg sync.WaitGroup out := make(chan int) for _, c := range cs { // 每来源一个搬运 G (fan-out 点多读者) wg.Add(1) go func(c <-chan int) { defer wg.Done() for v := range c { out <- v } }(c) } go func() { wg.Wait(); close(out) }() // 关键: 等全部写端退出才关
g.Go 的函数返回非 nil error 时派生 ctx 被 cancel; g.Wait() 只返回第一个错误, 其余被丢弃 —— 要全部错误就自己往切片里记。g, ctx := errgroup.WithContext(ctx) g.Go(func() error { return fetch(ctx) }) // 首个 err g.Go(func() error { return nil }) err := g.Wait() // → 只返回第一个 error; 派生 ctx 已被 cancel
g.Go 调用本身阻塞等令牌, 一行替掉手搓 sem chan + defer 归还。g, ctx := errgroup.WithContext(ctx) g.SetLimit(16) // 关键: Go 1.20+ 内置信号量 for _, t := range tasks { g.Go(func() error { return process(ctx, t) }) } // 满员时 g.Go 阻塞等令牌 — 这就是限流, 不是 bug
make(chan struct{}, n), 提交前 sem <- struct{}{} 拿额度, defer func() { <-sem }() 归还 (defer 保证 panic 路径也不漏)。sem := make(chan struct{}, 10) sem <- struct{}{} // 拿额度: 满则阻塞等 defer func() { <-sem }() // 关键: defer 归还, panic 路径不漏
close(done) 是广播, 唤醒所有 <-done; 它就是标准库 context 的思想祖宗。done := make(chan struct{}) close(done) // 关键: 关闭 = 广播, 唤醒所有 <-done <-done // → 立即返回, N 个监听者同时收到
select { case v := <-ch: use(v) default: // 关键: chan 空则立即走这 — 非阻塞 try drop() // 丢弃降级; 别放裸 for 里变忙轮询 }
out := make(chan string, 512) // 小缓冲 for line := range src { out <- line // 满则阻塞: 压力反向传导到源头 } // 下游写得慢 → 读端自动慢 — 天然限流不丢数据
index, 汇总端按下标归位 (res[it.i] = it.r), 顺序即可还原, 零锁零竞争。type indexed struct { i int; r Result } res := make([]Result, len(ids)) for it := range out { res[it.i] = it.r // 关键: 按下标归位, 顺序还原零锁 }
ctx.Done() 不是扭头就 return —— 先排空队列 (drain) 再退出, 存量任务不丢; 收尾处理用新的 Background ctx。select { case <-ctx.Done(): return drainAll(tasks) // 关键: 先排空再退, 存量不丢 case t, ok := <-tasks: if !ok { return nil } // 上游 close: 正常收工 handle(t) }
atomic.AddInt64 / LoadInt64 或 mutex; 普通 int++ 是数据竞争, go test -race 必炸。var done int64 go func() { process(); atomic.AddInt64(&done, 1) }() go func() { process(); atomic.AddInt64(&done, 1) }() fmt.Println(atomic.LoadInt64(&done)) // → 2; int++ 是竞态
50 万 URL 抓取, 目标站单 IP 限并发, 且任一 worker 拿到封禁错误就该全队撤退:
func crawl(ctx context.Context, urls []string) error { g, ctx := errgroup.WithContext(ctx) // 控制线: 首错即取消 jobs := make(chan string, 200) // 缓冲吸收突发, 不打爆内存 worker := func() error { for u := range jobs { if err := fetch(ctx, u); err != nil { return err } // 级联取消 } return nil } for i := 0; i < 8; i++ { g.Go(worker) } // 8 = 目标站单 IP 上限 feed: for _, u := range urls { select { case jobs <- u: case <-ctx.Done(): break feed // 某 worker 已挂: 停止投喂 } } close(jobs) // 发送方负责关闭 return g.Wait() // 只返回第一个错误 }
对比裸开 goroutine-per-URL: 峰值并发可控, 被封禁时 0.1s 内全队停止, 不再浪费请求配额。
批量查 2000 个订单状态, 下游只扛 8 并发, 但调用方要求返回顺序与入参一致:
func fetchAll(ctx context.Context, ids []string) []Result { sem := make(chan struct{}, 8) // 下游接口只扛 8 并发 out := make(chan indexed, len(ids)) // 容量=任务数: 写端永不阻塞 var wg sync.WaitGroup for i, id := range ids { wg.Add(1) go func(i int, id string) { // 1.22 前必须传参捕获 defer wg.Done() sem <- struct{}{} // 拿到令牌才出发 defer func() { <-sem }() // panic 路径也要归还 out <- indexed{i, callAPI(ctx, id)} // 带下标: 保序关键 }(i, id) } go func() { wg.Wait(); close(out) }() // 等全部写端退出才能关 res := make([]Result, len(ids)) for it := range out { res[it.i] = it.r } // 按下标归位还原顺序 return res }
2000 单从串行 4 分钟压到 8 并发 32 秒, 且各写各的下标 —— 无锁也无数据竞争。
全量读进内存必 OOM, 三阶段各一个 goroutine, 数据像水流过去, 任意时刻只在内存里驻留一小段:
func etl(ctx context.Context, path string) error { return writeStage(ctx, parseStage(ctx, readStage(ctx, path))) } // parse/write 同构: range 上游 func readStage(ctx context.Context, p string) <-chan string { out := make(chan string, 512) // 小缓冲: 背压自然形成 go func() { defer close(out) // 谁开 chan 谁收尾 f, err := os.Open(p) if err != nil { log.Print(err); return } defer f.Close() sc := bufio.NewScanner(f) for sc.Scan() { select { case out <- sc.Text(): // 满则阻塞: 反压上游停读 case <-ctx.Done(): return // 下游不要了就停读 } } }() return out }
512MB 容器跑通 10GB 文件, RSS 稳定在 90MB: 库写得慢, 读端就自动慢下来, 这就是背压。
老代码里 sem chan + defer 归去归来容易漏, SetLimit 一行等价且语义更直白:
func batchProcess(ctx context.Context, items []Task) error { g, ctx := errgroup.WithContext(ctx) g.SetLimit(16) // 内置信号量: 最多 16 并发 for _, t := range items { g.Go(func() error { // 拿不到令牌时阻塞在这 if err := process(ctx, t); err != nil { return fmt.Errorf("task %d: %w", t.ID, err) } return nil }) } return g.Wait() // 一行替掉手搓 sem chan + defer 归还 }
注意 g.Go 在满员时会阻塞提交循环 —— 这正是限流语义, 不是 bug。
金融对账 10 万账户: 单户失败不该连坐取消全局, 但要记进报表第二天人工复核:
func nightlyReconcile(ctx context.Context, accts []string) (Report, error) { var mu sync.Mutex ok, failed := 0, 0 // 部分结果也要记账 g, ctx := errgroup.WithContext(ctx) g.SetLimit(32) for _, acc := range accts { acc := acc g.Go(func() error { if err := reconcile(ctx, acc); err != nil { mu.Lock(); failed++; mu.Unlock() return nil // 不连坐: 记录后继续 } mu.Lock(); ok++; mu.Unlock() return nil }) } err := g.Wait() // 全部跑完才出报表 return Report{OK: ok, Failed: failed}, err }
若业务要求 fail-fast (一个错全停), 把 return nil 换成 return err 即可 —— 两种策略只差这一行。
面试与源码都绕不开: close(done) 广播 vs send 只唤醒一个, 这是 context 的思想原型:
// done 是 ctx 出现前的取消原语, 标准库 ctx 同款思想 func worker(done <-chan struct{}, jobs <-chan Job) { for { select { case j, ok := <-jobs: // ok=false: 上游 close if !ok { return } process(j) case <-done: // close(done) 是广播 return // 所有监听者同时收到 } } } quit := make(chan struct{}) close(quit) // 广播取消: 绝不是 quit <- struct{}{} (只唤醒一个)
不限并发时 500 张图同时解码, 300MB×500 直接 OOMKilled; 信号量把"在制品数量"钉死:
func processImages(ctx context.Context, keys []string) error { sem := make(chan struct{}, 10) // 10×30MB ≈ 内存封顶 var wg sync.WaitGroup for _, k := range keys { wg.Add(1) go func(k string) { defer wg.Done() select { case sem <- struct{}{}: // 先拿内存额度 case <-ctx.Done(): return } defer func() { <-sem }() // defer: panic 也归还 img, err := loadFromOSS(ctx, k) // 约 30MB 解码缓冲 if err != nil { return } uploadThumbnail(ctx, img, k) // 加工完即可被 GC }(k) } wg.Wait() return nil }
控并发就是控内存: 波峰 RSS 从无界的 15GB 变成稳定 300MB+基础开销。
K8s 发 SIGTERM 时正在排队的任务不能丢, 但也不能无限等 —— 排空 + 上限两件事都做:
func serve(ctx context.Context, tasks <-chan Task) error { for { select { case <-ctx.Done(): // 收到停止信号 return drainAll(tasks) // 先排空存量: 不丢任务 case t, ok := <-tasks: if !ok { return nil } // 上游 close: 正常收工 if err := handle(ctx, t); err != nil { return err } } } } func drainAll(tasks <-chan Task) error { bg := context.Background() // 收尾用新 ctx: 原已取消 for t := range tasks { handle(bg, t) } return nil }
排空也要配 deadline (外层 WithTimeout 包一层), 防止上游永不 close 把优雅退出拖成永久挂起。
跑批 40 分钟没任何输出会被值班当成挂了; worker 里直接 log 又太吵 —— 独立协程定时打快照:
var done int64 // 原子计数: int++ 是竞态 total := int64(len(tasks)) var wg sync.WaitGroup for _, t := range tasks { wg.Add(1) go func(t Task) { defer wg.Done() process(t); atomic.AddInt64(&done, 1) }(t) } go func() { // 独立进度协程: 只读快照 tk := time.NewTicker(5 * time.Second) defer tk.Stop() // Ticker 不 Stop 就泄漏 for range tk.C { p := atomic.LoadInt64(&done) log.Printf("progress %d/%d", p, total) if p == total { return } } }() wg.Wait()
RSS 缓涨 + goroutine 只增不减, 先抓全量栈再对代码, 定位到"发送端退出却没人 close":
// 症状: 上线 3 天 goroutine 2k → 180k, RSS 缓涨 // curl :6060/debug/pprof/goroutine?debug=1 抓栈定位: // 180000 @ ... main.(*Client).Watch watch.go:42 // 42: for msg := range c.stream // 发送端退出却没人 close func (c *Client) Watch(ctx context.Context) { defer c.wg.Done() for { select { // 修复: 阻塞点全挂 ctx case msg, ok := <-c.stream: if !ok { return } c.handle(ctx, msg) case <-ctx.Done(): // 断连/超时也能走人 return } } }
修复后 goroutine 曲线变成锯齿 (涨了就落), RSS 回落到基线; 原则是每个阻塞点都有退出分支。
go func() { wg.Wait(); close(out) }(), 所有写端退出后才轮到 close。// 错: merge 先 close(out) 再搬运 — 慢写端 send → panic go func() { wg.Wait(); close(out) }() // 对: 等全部写端退出后才轮到 close
// 错: close(ch) 之后又 ch <- v // → panic: send on closed channel, 整个进程崩 // 对: 关闭权只归发送方; 多发送者由协调者统一关
<-ctx.Done(); 上线盯 goroutine 数告警。// 错: for { process(<-ch) } — 阻塞点无退出分支 select { case v := <-ch: process(v) case <-ctx.Done(): return // 对: 每个阻塞点挂 Done }
g.Wait() 才是收口点。正解: return 前必 Wait, 且 error 要接住往上抛。for _, u := range urls { g.Go(work) } return nil // 错: 函数返回了, 协程还在飞 return g.Wait() // 对: Wait 是收口点, error 上抛
make(chan T) 后直接 send, 生产者永久阻塞, 表现为接口 hang 死。原因: 无缓冲是会合语义, 双方到齐才走。正解: 消费者先启, 或改有缓冲, 死锁时 pprof goroutine 一抓一个准。ch := make(chan int) ch <- 1 // 错: 无消费者 → 永久阻塞, 接口 hang 死 ch2 := make(chan int, 1) ch2 <- 1 // 对: 给缓冲 (或消费者先启) → 立即返回
// 错: for i := 0; i < 500; i++ { go worker() } — P99 涨 10 倍 for i := 0; i < 8; i++ { go worker() } // 对: worker 数对齐下游容量(连接数×0.8 起步), 非核数
for range ch 只有 close 才退出, 生产者忘了关 = 消费者永久挂着。正解: 设计期先回答"谁最后 close", 写进代码注释; 顺带用 ctx.Done 兜底退出。// 错: for v := range ch { … } 而生产者永不 close // → 消费者永久挂着 (pprof: chan receive) case <-ctx.Done(): return // 对: Done 兜底退出
// 错: 以为 case 写在第一个就有优先级 — select 均匀随机 select { // 对: 嵌套 — 先非阻塞试 Done case <-ctx.Done(): return default: } select { case v := <-data: use(v) } // 再整体阻塞
-race 必报 data race。原因: append 扩容与写指针都不是原子的。正解: 各写各的下标 res[i] = r, 或各建各的局部 slice 最后合并。// 错: 多个 G 同时 append(results, r) → -race 必报, // 扩容与写指针非原子, 随机丢条目 res := make([]Result, n) res[i] = r // 对: 各写各下标, 零竞争
select { case out <- v: case <-ctx.Done(): return }。out <- v // 错: 下游早退 → 永久阻塞在无人读的 chan select { // 对: 每个发送都带退出分支 case out <- v: case <-ctx.Done(): return }
sem <- struct{}{}。原因: panic 或提前 return 路径漏了 <-sem。正解: 拿到额度立刻 defer func() { <-sem }(), defer 是唯一保证。// 错: sem <- struct{}{} 后 panic/提前 return 漏了 <-sem // → 几小时后全部 G 卡在 sem <- struct{}{} sem <- struct{}{} defer func() { <-sem }() // 对: 拿到额度立刻 defer 归还
g.Go(func() error { defer func(){ _ = recover() }(); ... }) 或在框架层统一兜住转 error。// 错: g.Go(func() error { return work() }) — 不 recover, // worker 一个 nil 指针解引用 → 整服务崩 g.Go(func() (err error) { defer func() { if r := recover(); r != nil { err = fmt.Errorf("panic: %v", r) // 对: 自己兜 } }() return work() })
NewTimer + Reset 复用一个, 或升 Go 1.23+。// 错: for { select { case <-time.After(time.Second): … } } // 1.23 前每轮新 timer, 到期前不可回收 → 内存缓涨 t := time.NewTimer(time.Second) // 对: 循环外建一个 defer t.Stop() // 每轮用前 t.Reset(time.Second) 复用
go func(){ wg.Add(1); ... }, 主流程可能先 Wait 完直接返回。原因: Add 与 Wait 竞态。正解: 先 Add 后 go, Add 永远在启动协程的那一行之前。// 错: go func(){ wg.Add(1); … }() — Add/Wait 竞态, // 主流程可能先 Wait 完直接返回 wg.Add(1) // 对: 先 Add 后 go go func() { defer wg.Done(); work() }()
ch := make(chan Job, 1000) // 错: cap=1000 以为扛 1000 并发 — 第 1001 个照样阻塞 for i := 0; i < 8; i++ { go worker(ch) } // 对: 并发度=worker 数/SetLimit, 缓冲只管突发吸收
index 下标, 汇总端归位; 或显式 sort。out <- r // 错: 结果无下标, 聚合顺序每次跑都不同 out <- indexed{i, r} // 对: 结果带 index res[it.i] = it.r // → 汇总端按下标归位, 顺序还原
go func(u string){...}(u) 传参, 或循环内 u := u。// 错 (1.22 前): for _, u := range urls { go func(){ fetch(u) }() } // → N 个 G 全拿到最后一个 url (共用一份变量) for _, u := range urls { go func(u string) { fetch(u) }(u) // 对: 传参捕获 }
// 错: case <-ctx.Done(): return — 存量 3000 条任务蒸发 case <-ctx.Done(): return drainAll(tasks) // 对: 先排空(配 deadline)再退
done <- struct{}{} 只唤醒一个监听者, 其余 worker 永远等不到。正解: 广播取消用 close(done), 关闭语义是"唤醒全部"。done <- struct{}{} // 错: 只唤醒一个监听者 close(done) // 对: 广播 — 所有 <-done 同时就绪
results <- r, wg.Wait 永不返回, 整体死锁。原因: 无缓冲需要消费者同时在场。正解: results 给足缓冲 (≥任务数), 或消费与 Wait 并行 (go func(){ wg.Wait(); close(results) }() 先行)。results := make(chan R) wg.Wait() // 错: 无缓冲+先 Wait — worker 卡 results <- r, 死锁 // 对: 缓冲≥任务数, 或消费与 Wait 并行: results2 := make(chan R, n) go func() { wg.Wait(); close(results2) }()