10.2 channel 与 select
10.1 节的扫描器把两件事揉在了一个 goroutine 里:发现到期任务、执行关闭动作。这在小例子里看不出问题,但真实项目里「发现」和「处理」的节奏往往不同——扫描是定时的、轻量的,处理可能要写数据库、要重试、要限速。把它们拆开的标准做法就是一条队列。
Go 里队列的原生形态是 channel。它不只是一个线程安全队列,更是一套「谁在等谁」的通信协议:发送方阻塞、接收方阻塞、关闭广播,这些语义组合起来能表达很多同步意图。
本节把 TaskAPI 推进到:新增一条
jobschannel 作为到期任务队列,扫描器只负责投递,由消费者统一执行关闭动作。
10.2.1 channel 的本质:带类型的管道
声明方式是 make(chan 元素类型, 缓冲容量),例如 make(chan int) 是容量 0 的无缓冲 channel,make(chan int, 10) 是容量 10 的有缓冲 channel。
发送用 ch <- v,接收用 v := <-ch,方向都是「箭头指向」的语义:ch <- v 是往 channel 里塞,<-ch 是从里面取。
一个关键事实:channel 是引用类型。把它传给函数传的是同一个底层结构,不会复制队列内容。所以「往函数里传 channel 当队列」是自然写法。
10.2.2 无缓冲 vs 有缓冲
这是最需要建立直觉的一点:
| 类型 | 发送方行为 | 语义 |
|---|---|---|
无缓冲 make(chan T) | 阻塞,直到有接收方接手 | 一次「交接」,双方必须同时在场 |
有缓冲 make(chan T, n) | 缓冲未满时立即返回 | 一次「投递」,解耦双方节奏 |
ch := make(chan int, 2)
ch <- 1
ch <- 2
fmt.Println("缓冲写入完成, len =", len(ch), "cap =", cap(ch))
实测输出:
缓冲写入完成, len = 2 cap = 2
注意 len(ch) 返回的是当前缓冲里的元素个数,不是容量;容量用 cap(ch)。这两个值都只能当参考——在并发环境下读到的瞬间就可能过期,不要用它做逻辑判断。
无缓冲 channel 的交接语义:unbuf <- "任务队列" 这条发送语句会一直阻塞,直到另一个 goroutine 执行到 <-unbuf。这就是「同步点」——无缓冲 channel 天然是一把「双方碰面才放行」的闸门。
10.2.3 close 的语义:能读、不能再写
close(ch) 表示「不会再有新值了」。它的行为很容易记错,逐条列出:
| 操作 | 已关闭的 channel |
|---|---|
| 接收缓冲里剩余的值 | 正常返回剩余值 |
| 缓冲读空后再接收 | 立即返回零值,且第二个返回值 ok 为 false |
| 再发送 | panic:send on closed channel |
| 再次 close | panic:close of closed channel |
实测「读空后返回零值」:
ch := make(chan int, 2)
ch <- 10
ch <- 20
close(ch)
fmt.Println(<-ch) // 10
fmt.Println(<-ch) // 20
v, ok := <-ch
fmt.Println(v, ok) // 0 false
实测输出:
10
20
0 false
所以「接收」永远不会 panic,它用 ok == false 告诉你 channel 已经枯竭。这也是 for v := range ch 能自动结束的原因:range 内部就在检查 ok,一旦为 false 就退出循环。
反过来,「发送」到已关闭 channel 会 panic:
ch := make(chan int, 1)
close(ch)
ch <- 1 // panic: send on closed channel
实测 recover 到的信息:
recover 到 panic: send on closed channel
由此得到一条必须背下来的规则:只有发送方关闭 channel,接收方永远不要关。因为发送方知道「没有更多数据了」,接收方不知道。如果多个发送方,那要么由最后一个发送方关(需要额外同步),要么单独起一个协调 goroutine 在所有人都结束后关——10.3 节会用这个模式。
10.2.4 只发 / 只收类型
channel 可以声明方向,用来在类型层面约束函数职责:func produce(out chan<- int) 里的 chan<- int 表示「只能发」,func consume(in <-chan int) 里的 <-chan int 表示「只能收」。
方向类型可以隐式转换:chan int 能赋给 chan<- int 和 <-chan int,反过来不行。
这不是语法洁癖,而是把「谁负责 close」变成编译期约束的好办法:consume 拿到的是 <-chan int,它在函数体里根本写不出 close(in),编译器直接报错。用方向类型写签名,可以避免一大类「接收方误关 channel」的 bug。
10.2.5 select:同时等多路
select 是 channel 版的 switch:它在一组收发操作里挑一个「现在能执行」的来执行。
select {
case v := <-ch1:
fmt.Println("来自 ch1:", v)
case ch2 <- 42:
fmt.Println("成功写入 ch2")
default:
fmt.Println("一个都没就绪")
}
规则:
- 有多个 case 就绪时,随机选一个(避免饥饿)。
- 一个都没有且没有
default时,select阻塞,直到有 case 就绪。 - 有
default时,永不阻塞——一个都没就绪就走default。
第 1 点的随机性是刻意设计的,但它只在多个 case 同时就绪时才显现。下面这段代码里两个任务先后入队、stop 要等 30 毫秒后才触发,任一时刻只有一个 case 可执行,所以输出顺序其实是固定的:
tasks := make(chan string, 3)
stop := time.After(30 * time.Millisecond)
tasks <- "关闭到期任务 #1"
tasks <- "关闭到期任务 #2"
for {
select {
case t := <-tasks:
fmt.Println("处理:", t)
case <-stop:
fmt.Println("超时,退出")
return
}
}
实测输出:
处理: 关闭到期任务 #1
处理: 关闭到期任务 #2
超时,退出
time.After(d) 返回一个 <-chan time.Time,在 d 之后收到一个值。把它作为一个 case,就得到了「带超时的等待」——这是 select 最高频的用法之一。
10.2.6 三种常用 select 形态
把上面的规则组合起来,日常写代码基本就是这三种:
// 形态一:带超时地等一个结果
select {
case r := <-results:
handle(r)
case <-time.After(2 * time.Second):
return errors.New("处理超时")
}
// 形态二:非阻塞试探(典型用途是「尝试投递,投不进去就丢」)
select {
case jobs <- t:
// 入队成功
default:
log.Println("队列已满,丢弃任务", t.ID)
}
// 形态三:带退出信号的长循环(服务型 goroutine 的标准骨架)
for {
select {
case t := <-jobs:
process(t)
case <-done:
return
}
}
第 3 种你已经在 10.1 节见过一次。它是后面所有后台组件的基础模板:一个 for 循环 + 一个 select,一个 case 干活,一个 case 退出。
10.2.7 nil channel 会永久阻塞
零值的 channel 是 nil,对它的收发都会永久阻塞:
var nilCh chan int
select {
case <-nilCh:
fmt.Println("不可能走到这里")
default:
fmt.Println("nil channel 永远不可就绪")
}
实测输出:
nil channel 永远不可就绪
这个特性常被当作 bug,但它其实是个有用的工具:把某个 case 的 channel 设成 nil,就等价于关闭这一路,而不需要额外加布尔标志。例如「生产者结束后只想等退出信号,不想再收任务」:
jobsCh := jobs
for {
select {
case t, ok := <-jobsCh:
if !ok {
jobsCh = nil // 关掉这一路,后续 select 只会等 done
continue
}
process(t)
case <-done:
return
}
}
这个技巧后面讲优雅退出时会再用到,先记住结论:nil channel 上的收发永远阻塞,所以它在 select 里等于「永久不就绪」。
10.2.8 把任务队列接进 TaskAPI
现在把本节的东西组装起来:扫描器只投递,消费者负责关闭。生产者用 select 保证在收到停止信号时不会卡死,消费者用 ok 判断队列枯竭。
package main
import (
"fmt"
"time"
)
type Task struct {
ID int64
Title string
}
func main() {
jobs := make(chan Task, 5)
stop := make(chan struct{})
go func() {
defer close(jobs) // 谁发送谁关闭
for i := int64(1); i <= 4; i++ {
select {
case jobs <- Task{ID: i, Title: fmt.Sprintf("到期任务%d", i)}:
fmt.Printf("已投递 #%d\n", i)
case <-stop:
fmt.Println("生产者收到停止信号")
return
}
time.Sleep(10 * time.Millisecond)
}
}()
for {
select {
case t, ok := <-jobs:
if !ok {
fmt.Println("队列已关闭,消费者退出")
return
}
fmt.Printf("消费者关闭任务 #%d(%s)\n", t.ID, t.Title)
case <-time.After(100 * time.Millisecond):
fmt.Println("空闲超时,退出")
return
}
}
}
实测输出:
已投递 #1
消费者关闭任务 #1(到期任务1)
已投递 #2
消费者关闭任务 #2(到期任务2)
已投递 #3
消费者关闭任务 #3(到期任务3)
已投递 #4
消费者关闭任务 #4(到期任务4)
队列已关闭,消费者退出
这份代码把 10.1 节的两个职责真正拆开了:生产者只管投递(并且能被 stop 中断),消费者只管处理(并且能在队列枯竭时干净退出)。defer close(jobs) 保证了无论生产者从哪条路径返回,队列都会被关闭,消费者不会永远卡在 <-jobs 上。
10.2.9 常见错误清单
| 错误 | 后果 | 正确做法 |
|---|---|---|
接收方调用 close | 发送方 panic | 只让发送方关 |
重复 close | panic | 用 sync.Once 或明确单一关闭点 |
| 忘记关 channel | 消费者永久阻塞(泄漏) | 用 defer close |
用 len(ch) 判断是否为空 | 逻辑竞态 | 用 ok 或单独的信号 |
| 向 nil channel 发送 | 永久阻塞 | 初始化后再用 |
热循环里反复 time.After | 定时器堆积 | 复用 time.Timer |
10.2.10 小结与练习
- 无缓冲 channel 是「双方碰面」的同步点,有缓冲 channel 是「投递即可」的解耦器。
- 关闭的 channel 还能读空缓冲,读空后返回零值 +
ok=false;再发送会 panic。 - 发送方关闭是铁律,方向类型能帮编译器帮你守住它。
select在多路就绪时随机选;加default变非阻塞;nil channel 等于永久不就绪。
练习:
- 把 10.2.8 的消费者从 1 个改成 3 个(各自起一个 goroutine),观察输出顺序,并思考「谁该负责 close 结果 channel」。
- 用
default给生产者加一个「队列满就丢弃并计数」的分支,用atomic.Int64统计丢弃数。 - 把
stop换成一个context.Context的Done()channel,体会两者的关系(第 12 章会正式讲)。
下一节我们把单个消费者扩展成固定数量的 worker pool,并把「扫描 → 排队 → 处理 → 汇总」串成一条 pipeline。
阅读导航:上一节:10.1 goroutine 与调度直觉 · 下一节:10.3 worker pool 与 pipeline 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。