Go · Channel 与 select

CSP 模型: "不要通过共享内存来通信, 而要通过通信来共享内存" — hchan 结构 / close 语义 / select 多路复用

直接交接 满则排队 channel 底层 — runtime.hchan 自带锁的线程安全队列 (CSP 信道) buf — 环形缓冲队列 (有缓冲时) v v · · len()/cap() 就是它 sendq — 阻塞的发送 G 链表 buf 满时: 发送者挂起排队 (FIFO) recvq — 阻塞的接收 G 链表 buf 空时: 接收者挂起排队 互斥锁 lock 保护 + 直接从发送者栈拷贝到接收者栈 (无缓冲时) 无缓冲 chan — 同步交接 make(chan T) 发送者阻塞, 直到接收者出现 数据不进 buf, 直接栈到栈拷贝 会合点: 双方到齐才一起走 有缓冲 chan — 异步解耦 make(chan T, N) buf 未满: send 立即返回 buf 满: 发送 G 挂到 sendq 睡眠 生产者消费者解耦的关键 close 语义 — 必背 读已关闭 → 零值 + ok=false 向已关闭发送 → panic 重复 close → panic nil chan → 永久阻塞 close 是广播: 唤醒全部接收者 select — 多路复用 多个 chan 同时等, 谁就绪走谁 多个就绪 → 随机选 (防饿死) default → 无就绪立即走 超时 / 取消 / 扇入全靠它 四大惯用法 — 生产代码的标准形状 扇入 fan-in N 个来源的结果 汇入一条聚合 channel 聚合接口的标配 worker pool N 个 worker 消费 jobs chan worker 数 = 并发上限 保护下游/控制内存 timeout 熔断 select + time.After(d) 慢操作强制超时返回 每个下游调用都该有 ctx 取消传播 select 监听 ctx.Done() 用户断开 → 全链路停止 Go 服务的生命线 Legend channel 本体 (hchan) 无缓冲 / 惯用法 有缓冲 危险语义 select

本质: 带锁的安全队列

  • • hchan = 环形 buf + sendq/recvq + 锁
  • • 阻塞的是 goroutine 而非线程 — G 挂起, M 去干别的
  • • 无缓冲时数据栈到栈直拷, 零中转

缓冲 = 解耦程度

  • • cap 0: 强同步 (会合语义), 常用于信号
  • • cap 1: 单槽信号量 / done 标志
  • • cap N: 吸收突发, 但越大问题暴露越晚

select 是调度中枢

  • • 超时、取消、扇入、非阻塞读全靠 select
  • • 随机分支保证多就绪时不饿死
  • • nil 分支 = 动态禁用某个 case

💡 一句话理解

channel 是 Go 给并发通信开的"官方高速路": 一个自带锁的线程安全队列。发送/接收会阻塞goroutine而不是线程(阻塞的 G 挂起, M 继续跑别的), 所以 channel 编程模型虽"同步", 成本却极低。close 是广播信号, select 是多路分派 —— 两者组合出超时、取消、扇入扇出全部生产模式。

🧠 必知必会 必考 & 必会

谁应该 close
只有发送方 close, 且一个 channel 只 close 一次。多发送者时用额外的 done chan 或 sync.WaitGroup 协调, 由"最后一个发送者"或独立的 coordinator 关闭。
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 后自动退出
close 的广播语义
关闭后: 所有阻塞的接收者被唤醒, 之后每次读都立刻返回零值 + 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 自动退出
nil channel 语义
发送/接收都永久阻塞。看似无用, 实际是 select 里动态禁用分支的惯用法: 把不想参与的 case 置 nil, 该分支永远不会被选中。
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{}
}
hchan 细节
buf 是环形队列; 无缓冲发送时若 recvq 有人, 数据直接从发送者栈拷贝到接收者栈(零 buf 中转); 等待队列 FIFO, 保证公平。
ch := make(chan int, 3)     // buf 是环形队列
ch <- 1; ch <- 2
len(ch)                         // → 2 (buf 占用)
cap(ch)                         // → 3; 关键: 无缓冲+对端在等 → 栈到栈直拷
select 公平性
多个 case 同时就绪时均匀随机选择, 防止固定顺序饿死后面的 case; 这是语言规范保证, 不是实现巧合。
for i := 0; i < 10; i++ {
    select {
    case <-a:                    // a、b 同时就绪
    case <-b:
    }
}
// → a、b 各被选中约 5 次; 关键: 均匀随机是规范保证
与锁的分工
"传递数据所有权/事件通知"用 channel; "保护共享内存的复合不变量"用 mutex。两者不是对立, 是分工 —— 强行二选一是初学者宗教战争。
ch := make(chan Job)         // 传递"所有权": 给你, 我不碰
go func() { for j := range ch { run(j) } }()
ch <- job
var mu sync.Mutex             // 保护"不变量": 一起改才一致
mu.Lock(); n++; mu.Unlock()

🏭 生产实战 real world

场景 1 · worker pool 控并发保护下游

批量导出 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 优雅结束

场景 2 · 每个下游调用都套超时

select {
case v := <-queryCh:
    return v
case <-time.After(800 * time.Millisecond):
    return ErrTimeout                   // 熔断慢查询, 保住整体延迟
case <-ctx.Done():                         // 用户取消/整链超时
    return ctx.Err()
}

场景 3 · 扇入聚合多个数据源

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
}

场景 4 · 令牌桶限流 channel 化

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)                          // 调用点: 拿到令牌才放行

场景 5 · fan-out + 首错取消

并发打 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 }
}

场景 6 · 事件订阅: 广播与退出

发布订阅用 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: }     // 满则丢: 慢订阅者不拖垮发布
    }
}

场景 7 · 流水线 pipeline(生成器→阶段→汇)

经典三段式: 源生成、阶段变换、汇聚消费, 每阶段一个 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)) { ... }

场景 8 · 心跳通道: 判活与慢消费者检测

处理循环同时监听心跳与数据, 心跳超时说明上游卡死:

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
    }
}

场景 9 · 优雅停机三件套

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. 等存量请求完成, 超时强制断

场景 10 · 每请求超时贯穿到 DAO

超时不该只在 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)
}

⚠️ 编码注意与常见坑 pitfalls

坑 1 · 向已关闭的 channel 发送 — 直接 panic。正解: 发送方独占关闭权; 多发送者用 WaitGroup + 协调者收口, 或用 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 }()
坑 2 · goroutine 泄漏 — 接收者等一个永远不来的发送(或反之), goroutine 常驻 → 内存与调度压力缓涨。正解: 所有阻塞 select 必带 ctx.Done() 分支; 上线盯 pprof goroutine 数。
// 错: 结果无人收, 发送 G 永久卡死在 send
go func() { resultCh <- slowWork() }()
select {                          // 对: 阻塞处必带退出分支
case r := <-resultCh: use(r)
case <-ctx.Done():                // → 取消时发送方也应 select+ctx
}
坑 3 · 消费者忘记退出条件 — for range ch 只有 close 才结束; 生产者不 close = 消费者永远挂着。正解: 生命周期设计时先想清楚"谁最后 close"。
for job := range jobs {          // 只有 close(jobs) 才结束
    run(job)
}
// 错: 设计期没定"谁 close" → 消费者永挂
// 对: 唯一发送方 defer close(jobs), 责任写进注释
坑 4 · 循环里 time.After 泄漏(旧版本) — 1.23 前 time.After 创建的 timer 到期前不可回收, 高频循环里堆积。正解: 循环外复用 NewTimer + Reset, 或升 1.23+(timer 已可回收)。
for {                            // 错(1.22-): 每轮新建 timer 堆积
    select {
    case <-time.After(30 * time.Second): return
    case m := <-msg: handle(m)
    }
}
t := time.NewTimer(30 * time.Second)  // 对: 循环外复用+Reset
坑 5 · 用 len/cap 做流程控制 — len(ch) 是瞬时值, 拿到就过期, 用它判断"还有没有数据"必然竞态。正解: len 只用于监控指标; 流程控制用 close/ok/done 语义。
if len(ch) == 0 { /* 以为没数据 */ }   // 错: 瞬时值, 拿到就过期
select {                              // 对: 阻塞/非阻塞语义才可靠
case v := <-ch: use(v)
default: /* 此刻真无数据 */
}
坑 6 · 缓冲开太大掩盖背压 — cap 100000 = 问题延迟 10 万条才暴露, 内存先爆。正解: 缓冲 = 小缓冲(吸收突发)+ 监控队列长度告警, 而不是"大一点再大一点"。
jobs := make(chan Job, 100000)   // 错: 背压延迟 10 万条才暴露
jobs := make(chan Job, 64)         // 对: 小缓冲吸收突发
gauge("queue_depth", len(jobs))       // → 队列长度接告警
坑 7 · 生产者不 close, 消费者永等 — for range ch 只有 close 才退出。正解: 设计期就定"谁负责 close"(最后一个生产者/coordinator), 写进结构体注释。
// 错: 生产循环结束但忘了 close
for _, o := range orders { jobs <- o }
close(jobs)                      // 对: 最后一个生产动作后必关
for job := range jobs { run(job) }  // → 消费干净后正常退出
坑 8 · select 滥用 default 变忙轮询 — 带 default 的 select 不阻塞, 放进无 sleep 的 for 就是 100% CPU 空转。正解: 轮询必须配 time.Sleep/Ticker; 或去掉 default 走阻塞。
for {                            // 错: 无间隔 default = 100% CPU 空转
    select {
    case m := <-ch: use(m)
    default:
    }
}
for {                            // 对: 轮询配 Ticker, 或去掉 default
    select { case m := <-ch: use(m); case <-ticker.C: }
}
坑 9 · 多发送者场景下随手 close — 两个 goroutine 都往同一 chan 发, 任一方 close 都可能让对方 panic。正解: 外层 WaitGroup 收口后由唯一协调者 close, 或每人一条结果 chan 汇聚。
go func() { ch <- 1; close(ch) }()   // 错: 另一方可能 panic
go func() { ch <- 2 }()                // → send on closed channel
go func() { wg.Wait(); close(ch) }()     // 对: 协调者收口后唯一 close
坑 10 · 把缓冲容量当并发度 — cap=1000 的 chan 若无人消费, 第 1001 个发送照样阻塞; 缓冲只是"排队空间"不是"处理能力"。正解: 并发度 = worker 数, 缓冲 = 突发吸收, 两个概念分开调。
results := make(chan R, 1000)  // 错: 以为"能存=能处理"
results <- r                     // → 无人消费时第 1001 个照样阻塞
for i := 0; i < 20; i++ {          // 对: 并发度 = worker 数
    go worker(results)
}
坑 11 · channel 传大结构体全量拷贝 — 每次发送复制整个值(几百 KB 的结构体 = 明显开销)。正解: 传指针并明确"所有权已转移, 发送方不再碰"; 或定义好不可变约定。
ch1 <- frame                   // 错: chan Frame, 每次拷贝整个结构体
ch2 <- &frame                  // 对: chan *Frame, 只拷贝 8B 指针
// 约定: 发送后所有权归接收方, 发送方不再读写
坑 12 · 用 len(ch) 判断是否已关闭 — len 只反映当前积压, 关闭后残留数据 len 仍 >0。正解: 关闭状态用 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
坑 13 · 每轮循环新建 time.After/ctx — for-select 每轮 time.After(30s) 会创建新 timer(1.23 前不可回收, 高频循环堆积内存)。正解: 循环外 NewTimer+Reset, 或升 1.23+。
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)  // 对: 循环外建一次
坑 14 · chan of chan 的过度设计 — 动态注册/多级 chan 转发很快失控。正解: 大多数需求用 map + mutex 注册表 + 单一事件 chan 更直观。
reg := make(chan (chan Event))   // 错: chan 套 chan, 流转失控
mu.Lock(); subs[id] = make(chan Event, 16)  // 对: map+mutex 注册表
mu.Unlock()                        //    + 单一事件 chan 广播
坑 15 · 双向等待死锁 — A 等 B 的结果, B 等 A 的确认, 两个 chan 交错即永久互等。正解: 依赖图评审; 超时兜底(wait_for 语义用 ctx+select); 响应式设计避免握手。
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(): }  // 对: 超时兜底
坑 16 · 泄漏链条: goroutine 持 chan 持大对象 — 卡死的 goroutine 通过闭包引用着整个请求上下文, 内存泄漏被放大。正解: 取消机制 + 泄漏 goroutine 的内存账要在 pprof heap 里一起看。
// 错: 卡死的 G 闭包引用整个请求 ctx/bigResult → 泄漏放大
go func() { resultCh <- bigResult }()   // 无人收, 常驻
// 对: pprof goroutine+heap 一起看; 发送处 select+ctx.Done 退出
坑 17 · 关闭后读"以为没数据" — close 只代表"不会再有新数据", 缓冲里的残留值仍可依次读出(这是特性, 用于收尾)。正解: 收尾消费用 range 读干净再退出; 需要丢弃语义就重建 chan。
ch := make(chan int, 3)
ch <- 1; ch <- 2; close(ch)
for v := range ch { fmt.Print(v) }  // → 12, 残留照常读出
// 关键: close = "不会再有新数据", 不是"清空"
坑 18 · 依赖 select 分支的书写顺序 — 多个 case 就绪时是随机选择, 与代码顺序无关; 按顺序写的"优先级"是错觉。正解: 需要优先级就嵌套 select(先试高优先级+default, 再整体阻塞)。
select {
case <-high:                    // 错: 以为"排在前面"有优先级
case <-low:                     //    实际随机, 顺序是错觉
}
select {                        // 对: 真优先级要嵌套
case <-high:
default:
    select { case <-high: case <-low: }  // 高优无则一起等
}
坑 19 · 用 chan 手搓 mutex — cap=1 chan 当锁用理论可行, 但不可重入、难调试、易忘"解锁"。正解: 直接 sync.Mutex; chan 只用于数据/事件语义。
lock := make(chan struct{}, 1)
lock <- struct{}{}               // "加锁"
critical()
<-lock                           // 错: 忘了这行 = 永久"持有"
var mu sync.Mutex                // 对: 直接 Mutex, 语义清晰
mu.Lock(); critical(); mu.Unlock()
坑 20 · nil channel 语义用错方向 — nil chan 收发永久阻塞(可用于 select 禁用分支), 但把它当"空 channel"发数据 = 永久卡死且无报错。正解: nil 赋值只出现在"动态关闭某分支"的惯用法里; 意外的 nil 多半是初始化遗漏。
var ch chan int              // 错: 忘了 make, 当"空 chan"用
ch <- 1                         // → 永久阻塞, 无报错难排查
ch = make(chan int, 1); ch <- 1  // 对: 先 make
var off <-chan int          // 关键: nil 只用于禁用分支惯用法