CSP 模型: "不要通过共享内存来通信, 而要通过通信来共享内存" — hchan 结构 / close 语义 / select 多路复用
channel 是 Go 给并发通信开的"官方高速路": 一个自带锁的线程安全队列。发送/接收会阻塞goroutine而不是线程(阻塞的 G 挂起, M 继续跑别的), 所以 channel 编程模型虽"同步", 成本却极低。close 是广播信号, select 是多路分派 —— 两者组合出超时、取消、扇入扇出全部生产模式。
ch := make(chan int) go func() { defer close(ch) // 关键: 发送方 defer close, 只关一次 for i := 0; i < 3; i++ { ch <- i } }() for v := range ch { fmt.Print(v) } // → 012 后自动退出
v, ok := <-ch 的 ok=false; for range ch 优雅退出 —— 这是"生产结束通知消费"的标准方式。ch := make(chan int, 5); ch <- 7; close(ch) v, ok := <-ch // → 7, true (残留值先出) v, ok = <-ch // → 0, false (排干后 ok=false) for v := range ch2 { } // 关键: range 遇 close 自动退出
var ch chan int // 未 make = nil, 收发永久阻塞 select { case v := <-ch: // 关键: select 里 nil 分支被禁用 _ = v case <-time.After(time.Second): // → 1s 后走这里 }
chan<- T / <-chan T 是函数签名里的"意图文档": 消费函数拿只读 chan, 天然防止越权发送/关闭 —— 编译期挡住一类事故。func consume(ch <-chan Event) { // 只读: 编译期禁 send/close for e := range ch { handle(e) } } func produce(ch chan<- Event) { // 只写: 编译期禁 receive ch <- Event{} }
ch := make(chan int, 3) // buf 是环形队列 ch <- 1; ch <- 2 len(ch) // → 2 (buf 占用) cap(ch) // → 3; 关键: 无缓冲+对端在等 → 栈到栈直拷
for i := 0; i < 10; i++ { select { case <-a: // a、b 同时就绪 case <-b: } } // → a、b 各被选中约 5 次; 关键: 均匀随机是规范保证
ch := make(chan Job) // 传递"所有权": 给你, 我不碰 go func() { for j := range ch { run(j) } }() ch <- job var mu sync.Mutex // 保护"不变量": 一起改才一致 mu.Lock(); n++; mu.Unlock()
批量导出 10 万单, 下游 DB 只能扛 20 并发:
jobs := make(chan Order, 1000) // 有缓冲: 吸收突发, 反压生产者 results := make(chan Result, 1000) var wg sync.WaitGroup for i := 0; i < 20; i++ { // worker 数 = 下游能扛的并发 wg.Add(1) go func() { defer wg.Done() for order := range jobs { // close(jobs) 后自动退出 results <- process(order) } }() } go func() { for _, o := range allOrders { jobs <- o }; close(jobs) }() wg.Wait(); close(results) // 消费者 range results 优雅结束
select { case v := <-queryCh: return v case <-time.After(800 * time.Millisecond): return ErrTimeout // 熔断慢查询, 保住整体延迟 case <-ctx.Done(): // 用户取消/整链超时 return ctx.Err() }
func merge(ctx context.Context, chans ...<-chan Item) <-chan Item { out := make(chan Item) var wg sync.WaitGroup wg.Add(len(chans)) for _, c := range chans { go func(c <-chan Item) { // 每个来源一个搬运 goroutine defer wg.Done() for it := range c { select { case out <- it: case <-ctx.Done(): return } } }(c) } go func() { wg.Wait(); close(out) }() return out }
Ticker 定期投放令牌, 缓冲即容量; 拿不到令牌的请求在接收处排队:
func limiter(rate int, burst int) <-chan struct{} { tokens := make(chan struct{}, burst) // 缓冲 = 桶容量 go func() { t := time.NewTicker(time.Second / time.Duration(rate)) for range t.C { select { case tokens <- struct{}{}: default: } // 满则丢弃, 不阻塞泵 } }() return tokens } <-limiter(100, 50) // 调用点: 拿到令牌才放行
并发打 N 个请求, 第一个失败就取消其余, 不浪费下游配额:
func fetchAll(ctx context.Context, urls []string) ([]Result, error) { ctx, cancel := context.WithCancel(ctx) defer cancel() // 兜底: 正常返回也释放 errs := make(chan error, 1) // 缓冲1: 首个错误不阻塞报告者 results := make([]Result, len(urls)) var wg sync.WaitGroup for i, u := range urls { wg.Add(1) go func(i int, u string) { defer wg.Done() r, err := fetchOne(ctx, u) // ctx 取消时快速失败返回 if err != nil { select { case errs <- err: cancel(): default: } // 只报第一个错 return } results[i] = r // 各写各下标: 无竞争 }(i, u) } wg.Wait() select { case err := <-errs: return nil, err: default: return results, nil } }
发布订阅用 chan 切片 + close 广播; 订阅者 range 即自动退出:
type Bus struct { mu sync.Mutex subs map[chan Event]struct{} // 每订阅者一条 chan } func (b *Bus) Subscribe() <-chan Event { ch := make(chan Event, 16) // 小缓冲吸收突发 b.mu.Lock(); b.subs[ch] = struct{}{}; b.mu.Unlock() return ch } func (b *Bus) Publish(e Event) { b.mu.Lock(); defer b.mu.Unlock() for ch := range b.subs { select { case ch <- e: default: } // 满则丢: 慢订阅者不拖垮发布 } }
经典三段式: 源生成、阶段变换、汇聚消费, 每阶段一个 goroutine:
func gen(ctx context.Context, nums ...int) <-chan int { out := make(chan int) go func() { defer close(out) for _, n := range nums { select { case out <- n: case <-ctx.Done(): return } } }() return out } func sq(in <-chan int) <-chan int { // 阶段: 变换 out := make(chan int) go func() { defer close(out) for n := range in { out <- n * n } // 上游 close 后自动结束 }() return out } // 消费: for v := range sq(gen(ctx, 1, 2, 3)) { ... }
处理循环同时监听心跳与数据, 心跳超时说明上游卡死:
heartbeat := time.NewTicker(5 * time.Second) defer heartbeat.Stop() for { select { case item := <-dataCh: process(item) // 正常路径 case <-heartbeat.C: // 5s 没数据也醒来: 记活/上报 markAlive() case <-ctx.Done(): return } }
ctx 取消 + WaitGroup 等待 + done 通知, 三者配合不丢任务:
ctx, cancel := context.WithCancel(context.Background()) go handleSignals(cancel) // SIGTERM → cancel() <-ctx.Done() // 1. 收到停止信号 server.SetKeepAlivesEnabled(false) // 2. 不再接新连接 ctxShutdown, c2 := context.WithTimeout(context.Background(), 15*time.Second) defer c2() server.Shutdown(ctxShutdown) // 3. 等存量请求完成, 超时强制断
超时不该只在 HTTP 层 —— ctx 一路传下去, 慢 SQL 一起被掐断:
func (h *Handler) GetOrder(w http.ResponseWriter, r *http.Request) { ctx, cancel := context.WithTimeout(r.Context(), 800*time.Millisecond) defer cancel() order, err := h.dao.GetOrder(ctx, id) // dao 用 ctx 执行查询 if errors.Is(err, context.DeadlineExceeded) { http.Error(w, "timeout", http.StatusGatewayTimeout) return // 上层 800ms 处置, 不拖 worker } json.NewEncoder(w).Encode(order) }
sync.Once 关 done。ch := make(chan int, 1) close(ch); ch <- 1 // 错: → panic: send on closed channel // 对: 发送方独占关闭权, close 放发送方 defer 里 go func() { defer close(ch); ch <- 1 }()
ctx.Done() 分支; 上线盯 pprof goroutine 数。// 错: 结果无人收, 发送 G 永久卡死在 send go func() { resultCh <- slowWork() }() select { // 对: 阻塞处必带退出分支 case r := <-resultCh: use(r) case <-ctx.Done(): // → 取消时发送方也应 select+ctx }
for range ch 只有 close 才结束; 生产者不 close = 消费者永远挂着。正解: 生命周期设计时先想清楚"谁最后 close"。for job := range jobs { // 只有 close(jobs) 才结束 run(job) } // 错: 设计期没定"谁 close" → 消费者永挂 // 对: 唯一发送方 defer close(jobs), 责任写进注释
for { // 错(1.22-): 每轮新建 timer 堆积 select { case <-time.After(30 * time.Second): return case m := <-msg: handle(m) } } t := time.NewTimer(30 * time.Second) // 对: 循环外复用+Reset
if len(ch) == 0 { /* 以为没数据 */ } // 错: 瞬时值, 拿到就过期 select { // 对: 阻塞/非阻塞语义才可靠 case v := <-ch: use(v) default: /* 此刻真无数据 */ }
jobs := make(chan Job, 100000) // 错: 背压延迟 10 万条才暴露 jobs := make(chan Job, 64) // 对: 小缓冲吸收突发 gauge("queue_depth", len(jobs)) // → 队列长度接告警
for range ch 只有 close 才退出。正解: 设计期就定"谁负责 close"(最后一个生产者/coordinator), 写进结构体注释。// 错: 生产循环结束但忘了 close for _, o := range orders { jobs <- o } close(jobs) // 对: 最后一个生产动作后必关 for job := range jobs { run(job) } // → 消费干净后正常退出
for { // 错: 无间隔 default = 100% CPU 空转 select { case m := <-ch: use(m) default: } } for { // 对: 轮询配 Ticker, 或去掉 default select { case m := <-ch: use(m); case <-ticker.C: } }
go func() { ch <- 1; close(ch) }() // 错: 另一方可能 panic go func() { ch <- 2 }() // → send on closed channel go func() { wg.Wait(); close(ch) }() // 对: 协调者收口后唯一 close
results := make(chan R, 1000) // 错: 以为"能存=能处理" results <- r // → 无人消费时第 1001 个照样阻塞 for i := 0; i < 20; i++ { // 对: 并发度 = worker 数 go worker(results) }
ch1 <- frame // 错: chan Frame, 每次拷贝整个结构体 ch2 <- &frame // 对: chan *Frame, 只拷贝 8B 指针 // 约定: 发送后所有权归接收方, 发送方不再读写
v, ok := <-ch 或 range 的退出行为判断。ch := make(chan int, 1); ch <- 3; close(ch) fmt.Println(len(ch)) // → 1 (错判: 以为还开着) v, ok := <-ch // 对: → 3, true; 再读 → 0, false
for { // 错: 每轮新建 ctx/timer, 高频循环堆积 ctx, cancel := context.WithTimeout(base, time.Second) select { case <-ctx.Done(): cancel(); return case m := <-ch: handle(m) } cancel() } ctx, cancel := context.WithTimeout(base, 30*time.Second) // 对: 循环外建一次
reg := make(chan (chan Event)) // 错: chan 套 chan, 流转失控 mu.Lock(); subs[id] = make(chan Event, 16) // 对: map+mutex 注册表 mu.Unlock() // + 单一事件 chan 广播
a, b := make(chan int), make(chan int) go func() { b <- 2; <-a }() // 错: 双向握手互等 a <- 1; <-b // → fatal error: all goroutines are asleep select { case resp <- r: case <-ctx.Done(): } // 对: 超时兜底
// 错: 卡死的 G 闭包引用整个请求 ctx/bigResult → 泄漏放大 go func() { resultCh <- bigResult }() // 无人收, 常驻 // 对: pprof goroutine+heap 一起看; 发送处 select+ctx.Done 退出
ch := make(chan int, 3) ch <- 1; ch <- 2; close(ch) for v := range ch { fmt.Print(v) } // → 12, 残留照常读出 // 关键: close = "不会再有新数据", 不是"清空"
select { case <-high: // 错: 以为"排在前面"有优先级 case <-low: // 实际随机, 顺序是错觉 } select { // 对: 真优先级要嵌套 case <-high: default: select { case <-high: case <-low: } // 高优无则一起等 }
lock := make(chan struct{}, 1) lock <- struct{}{} // "加锁" critical() <-lock // 错: 忘了这行 = 永久"持有" var mu sync.Mutex // 对: 直接 Mutex, 语义清晰 mu.Lock(); critical(); mu.Unlock()
var ch chan int // 错: 忘了 make, 当"空 chan"用 ch <- 1 // → 永久阻塞, 无报错难排查 ch = make(chan int, 1); ch <- 1 // 对: 先 make var off <-chan int // 关键: nil 只用于禁用分支惯用法