Appearance
Channel 进阶:select 与多路复用
上一篇我们学习了 channel 的基本用法,但实际项目中往往需要同时处理多个 channel——监听多个数据源、实现超时控制、响应取消信号等。这时就需要 select 语句。select 是 Go 并发编程的瑞士军刀,它能让一个 goroutine 同时等待多个 channel 操作,实现多路复用、超时、非阻塞读写等关键模式。本篇我们将全面掌握 select 的用法。
一、select 语句:多路复用
select 语句的语法类似 switch,但它是专门为 channel 设计的。它会同时监听多个 case,每个 case 必须是一个 channel 操作(发送或接收)。当其中某个 case 可以执行时,就执行对应的分支。
1. 基本语法
go
select {
case v := <-ch1:
// 从 ch1 收到数据
case ch2 <- data:
// 向 ch2 发送成功
case <-time.After(time.Second):
// 超时
default:
// 所有 case 都不就绪时执行
}2. 基本示例:等待两个 channel
go
package main
import (
"fmt"
"time"
)
func main() {
ch1 := make(chan string)
ch2 := make(chan string)
go func() {
time.Sleep(1 * time.Second)
ch1 <- "来自 ch1"
}()
go func() {
time.Sleep(2 * time.Second)
ch2 <- "来自 ch2"
}()
// 用 select 等待第一个就绪的 channel
for i := 0; i < 2; i++ {
select {
case msg := <-ch1:
fmt.Println("收到:", msg)
case msg := <-ch2:
fmt.Println("收到:", msg)
}
}
fmt.Println("完成")
}运行后,1 秒时会先打印「来自 ch1」,2 秒时再打印「来自 ch2」。select 会阻塞,直到某个 case 就绪。
3. select 的关键特性
理解 select 的行为有几个关键点:
- 阻塞:如果没有
default且没有 case 就绪,select会阻塞,直到某个 case 就绪。 - 就绪:一个 case 就绪指的是它的 channel 操作可以立即完成(接收有数据、发送有空间、channel 已关闭)。
- 随机选择:如果多个 case 同时就绪,
select会随机选一个执行,而不是按顺序。 - 无 default 时阻塞:没有
default的select是阻塞式的。 - 有 default 时非阻塞:有
default的select在所有 case 都不就绪时立即执行default,不阻塞。
二、select 的随机选择机制
这是 select 最容易被误解的特性:当多个 case 同时就绪时,select 随机选一个,而不是选第一个。这是刻意的设计,为了防止某个 case 因位置靠后而饥饿。
go
package main
import "fmt"
func main() {
ch1 := make(chan int, 1)
ch2 := make(chan int, 1)
// 两个 channel 都有数据,同时就绪
ch1 <- 1
ch2 <- 2
count1, count2 := 0, 0
for i := 0; i < 10000; i++ {
// 重新填满
ch1 <- 1
ch2 <- 2
select {
case <-ch1:
count1++
case <-ch2:
count2++
}
}
fmt.Printf("ch1 被选中 %d 次, ch2 被选中 %d 次\n", count1, count2)
// 两者大致各占 50%
}运行多次你会发现,两个 case 被选中的次数大致相等(约 5000 次各)。这与 switch 按顺序匹配的行为完全不同。
为什么是随机的
如果按顺序匹配,那么排在前的 case 只要一直就绪,排在后的就永远没机会——这就是「饥饿」。随机选择保证了公平性,每个 case 都有同等机会被选中。
三、default 分支:非阻塞操作
default 分支在所有 case 都不就绪时执行,利用它可以实现非阻塞的 channel 操作。
1. 非阻塞接收
go
package main
import "fmt"
func main() {
ch := make(chan int, 1)
// 非阻塞接收:没数据也不卡住
select {
case v := <-ch:
fmt.Println("收到:", v)
default:
fmt.Println("没有数据,跳过")
}
// 放入数据再试
ch <- 42
select {
case v := <-ch:
fmt.Println("收到:", v)
default:
fmt.Println("没有数据,跳过")
}
}2. 非阻塞发送
go
package main
import "fmt"
func main() {
ch := make(chan int, 1)
ch <- 1 // 缓冲区满了
// 非阻塞发送:满了就丢弃
select {
case ch <- 2:
fmt.Println("发送成功")
default:
fmt.Println("缓冲区满,丢弃")
}
}3. 实现一个简单的轮询
go
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan int, 5)
go func() {
for i := 1; i <= 3; i++ {
ch <- i
time.Sleep(300 * time.Millisecond)
}
}()
for i := 0; i < 10; i++ {
select {
case v := <-ch:
fmt.Println("收到:", v)
default:
fmt.Println("暂时没数据,做点别的")
}
time.Sleep(200 * time.Millisecond)
}
}default 让 select 变成「试一下,不行就算了」,这在需要保持 goroutine 响应性的场景很有用。
四、超时控制:select + time.After
time.After(d) 返回一个 channel,在 d 时间后会发送一个时间值。配合 select,可以优雅地实现超时控制——这是 Go 并发编程中最常用的模式之一。
1. 基本超时
go
package main
import (
"fmt"
"time"
)
func slowOperation() <-chan string {
ch := make(chan string)
go func() {
time.Sleep(2 * time.Second) // 模拟慢操作
ch <- "结果"
}()
return ch
}
func main() {
select {
case res := <-slowOperation():
fmt.Println("收到结果:", res)
case <-time.After(1 * time.Second):
fmt.Println("超时!操作太慢")
}
}1 秒后超时分支被选中,主 goroutine 不再等待慢操作。
2. 注意 time.After 的资源问题
time.After 每次调用都会创建一个新的 timer,在超时前不会被 GC。如果在高频循环里大量使用,可能导致内存增长。性能敏感场景可以用 time.NewTimer 手动管理:
go
package main
import (
"fmt"
"time"
)
func main() {
timer := time.NewTimer(1 * time.Second)
defer timer.Stop() // 用完停止,及时释放
ch := make(chan string)
go func() {
time.Sleep(2 * time.Second)
ch <- "结果"
}()
select {
case res := <-ch:
fmt.Println("收到:", res)
case <-timer.C:
fmt.Println("超时")
}
}3. 多级超时
go
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan string)
go func() {
time.Sleep(500 * time.Millisecond)
ch <- "数据"
}()
select {
case v := <-ch:
fmt.Println("快速收到:", v)
case <-time.After(200 * time.Millisecond):
fmt.Println("第一级超时,再等等")
select {
case v := <-ch:
fmt.Println("终于收到:", v)
case <-time.After(1 * time.Second):
fmt.Println("彻底超时")
}
}
}五、tickers 和 timers
time 包提供了两个定时器工具,常与 select 配合使用。
1. Timer(定时器)
time.NewTimer(d) 创建一个只触发一次的定时器,到时间后向 timer.C 发送一个值:
go
package main
import (
"fmt"
"time"
)
func main() {
timer := time.NewTimer(3 * time.Second)
fmt.Println("等待 3 秒...")
<-timer.C // 阻塞直到定时器触发
fmt.Println("3 秒到了!")
// 可以提前 Stop
timer2 := time.NewTimer(5 * time.Second)
go func() {
<-timer2.C
fmt.Println("这个不会执行")
}()
timer2.Stop() // 立即停止
fmt.Println("timer2 已停止")
time.Sleep(time.Second)
}2. Ticker(周期触发器)
time.NewTicker(d) 创建一个周期性触发的定时器,每隔 d 时间向 ticker.C 发送一个值:
go
package main
import (
"fmt"
"time"
)
func main() {
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop() // 记得 Stop,否则 ticker 不会停止
done := make(chan struct{})
go func() {
time.Sleep(2500 * time.Millisecond)
close(done)
}()
count := 0
loop:
for {
select {
case <-ticker.C:
count++
fmt.Printf("tick %d\n", count)
case <-done:
fmt.Println("结束")
break loop
}
}
}3. 实际应用:限速器
go
package main
import (
"fmt"
"time"
)
func main() {
// 每 200ms 允许一次操作 = 每秒 5 次
ticker := time.NewTicker(200 * time.Millisecond)
defer ticker.Stop()
for i := 1; i <= 5; i++ {
<-ticker.C // 等待令牌
fmt.Printf("执行操作 %d, 时间: %v\n", i, time.Now().Format("15:04:05.000"))
}
fmt.Println("完成")
}六、用 select 实现非阻塞读写
结合 select 和 default,可以实现各种非阻塞操作。
1. 非阻塞遍历 channel
go
package main
import "fmt"
func main() {
ch := make(chan int, 5)
for i := 1; i <= 3; i++ {
ch <- i
}
// 把 channel 里的数据一次性取空,取完不阻塞
for {
select {
case v := <-ch:
fmt.Println("取出:", v)
default:
fmt.Println("channel 空了")
return
}
}
}2. 优先级 channel
用嵌套 select 实现「优先 channel」:先非阻塞试优先 channel,不行再等所有 channel:
go
package main
import (
"fmt"
"time"
)
func main() {
highPri := make(chan string, 1)
lowPri := make(chan string, 1)
// 先放低优先级,再放高优先级
lowPri <- "低优先级消息"
highPri <- "高优先级消息"
for i := 0; i < 2; i++ {
// 第一层:非阻塞试高优先级
select {
case v := <-highPri:
fmt.Println("处理:", v)
continue
default:
}
// 第二层:阻塞等待任意一个
select {
case v := <-highPri:
fmt.Println("处理:", v)
case v := <-lowPri:
fmt.Println("处理:", v)
}
}
time.Sleep(time.Second)
}七、select 中的 nil channel
这是一个高级技巧:对 nil channel 的操作会永久阻塞,但在 select 中,nil channel 的 case 永远不会被选中。利用这一点,可以动态「禁用」某个 case。
1. nil channel 的行为
go
package main
import "fmt"
func main() {
var ch chan int // nil channel
fmt.Println("ch 是 nil:", ch == nil)
// <-ch // 永久阻塞(死锁)
// ch <- 1 // 永久阻塞(死锁)
// close(ch) // panic
}2. 用 nil channel 动态启用/禁用 case
典型场景:某个 channel 还没准备好时,先设为 nil,select 会忽略它;准备好后赋值,select 开始监听。
go
package main
import (
"fmt"
"time"
)
func main() {
var ch1 chan string = nil // 初始禁用
ch2 := make(chan string)
go func() {
time.Sleep(500 * time.Millisecond)
ch2 <- "来自 ch2"
}()
go func() {
time.Sleep(1 * time.Second)
ch1 = make(chan string) // 1 秒后启用 ch1
ch1 <- "来自 ch1"
}()
// 先等 ch2,然后启用 ch1
for i := 0; i < 2; i++ {
select {
case v := <-ch1:
fmt.Println("收到:", v)
case v := <-ch2:
fmt.Println("收到:", v)
ch2 = nil // 收过一次后禁用 ch2
}
}
time.Sleep(time.Second)
}这种模式在状态机、动态管道中很有用。
八、退出信号模式:用 select 监听 done channel
这是 Go 并发中最经典的模式之一:用一个 done channel 作为退出信号,goroutine 通过 select 监听它,收到信号就退出。
1. 基本模式
go
package main
import (
"fmt"
"time"
)
func worker(done <-chan struct{}) {
for {
select {
case <-done:
fmt.Println("worker 收到退出信号,结束")
return
default:
fmt.Println("worker 工作中...")
time.Sleep(300 * time.Millisecond)
}
}
}
func main() {
done := make(chan struct{})
go worker(done)
time.Sleep(2 * time.Second)
close(done) // 通知 worker 退出
time.Sleep(500 * time.Millisecond)
fmt.Println("主 goroutine 结束")
}2. 结合数据 channel 和 done
go
package main
import (
"fmt"
"time"
)
func processor(input <-chan int, output chan<- int, done <-chan struct{}) {
for {
select {
case <-done:
fmt.Println("processor 退出")
return
case v, ok := <-input:
if !ok {
fmt.Println("input 关闭,processor 退出")
return
}
output <- v * 2
}
}
}
func main() {
input := make(chan int, 5)
output := make(chan int, 5)
done := make(chan struct{})
go processor(input, output, done)
// 投递任务
for i := 1; i <= 3; i++ {
input <- i
fmt.Println("产出:", <-output)
}
close(done) // 通知退出
time.Sleep(time.Second)
}这个模式是 context 包的基础——后面我们会看到,context.Done() 返回的就是一个这样的 done channel。
九、扇入模式:合并多个 channel
扇入(fan-in)模式把多个输入 channel 合并成一个输出 channel,是 select 的经典应用。
go
package main
import (
"fmt"
"sync"
"time"
)
// 把多个 channel 合并成一个
func fanIn(channels ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
// 为每个输入 channel 启动一个 goroutine
for _, ch := range channels {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for v := range c {
out <- v // 转发到输出
}
}(ch)
}
// 所有输入结束后关闭输出
go func() {
wg.Wait()
close(out)
}()
return out
}
func source(name string, start, count int) <-chan int {
ch := make(chan int)
go func() {
defer close(ch)
for i := start; i < start+count; i++ {
ch <- i
time.Sleep(200 * time.Millisecond)
}
fmt.Printf("%s 完成\n", name)
}()
return ch
}
func main() {
ch1 := source("A", 1, 3)
ch2 := source("B", 10, 3)
ch3 := source("C", 20, 3)
// 合并三个 channel
merged := fanIn(ch1, ch2, ch3)
for v := range merged {
fmt.Println("合并收到:", v)
}
fmt.Println("全部完成")
}扇入模式让消费者只需要监听一个 channel,无需关心数据来自哪里。
十、完整示例:多数据源并发查询
综合运用 select、超时、扇入,实现一个多数据源并发查询的例子:
go
package main
import (
"context"
"fmt"
"math/rand"
"sync"
"time"
)
type Result struct {
Source string
Data string
Err error
}
// 模拟从某个数据源查询
func query(ctx context.Context, source string, delay time.Duration) <-chan Result {
ch := make(chan Result, 1)
go func() {
defer close(ch)
select {
case <-time.After(delay):
// 模拟随机失败
if rand.Intn(10) < 2 {
ch <- Result{Source: source, Err: fmt.Errorf("%s 查询失败", source)}
} else {
ch <- Result{Source: source, Data: fmt.Sprintf("%s的数据", source)}
}
case <-ctx.Done():
ch <- Result{Source: source, Err: ctx.Err()}
}
}()
return ch
}
// 扇入多个查询结果
func fanInResults(ctx context.Context, sources ...<-chan Result) <-chan Result {
out := make(chan Result)
var wg sync.WaitGroup
for _, s := range sources {
wg.Add(1)
go func(c <-chan Result) {
defer wg.Done()
for r := range c {
select {
case out <- r:
case <-ctx.Done():
return
}
}
}(s)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
rand.Seed(time.Now().UnixNano())
// 设置总超时 2 秒
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
// 并发查询三个数据源
ch1 := query(ctx, "MySQL", 500*time.Millisecond)
ch2 := query(ctx, "Redis", 800*time.Millisecond)
ch3 := query(ctx, "ES", 1500*time.Millisecond)
// 扇入结果
results := fanInResults(ctx, ch1, ch2, ch3)
var success, failed int
for r := range results {
if r.Err != nil {
fmt.Printf("[失败] %s: %v\n", r.Source, r.Err)
failed++
} else {
fmt.Printf("[成功] %s: %s\n", r.Source, r.Data)
success++
}
}
fmt.Printf("\n总计: 成功 %d, 失败 %d\n", success, failed)
}这个例子展示了真实场景中的并发查询:多个数据源同时查,用 select + context 做超时控制,用扇入收集结果。这是微服务架构中「并发聚合」的常见模式。
十一、select 的几个细节
1. 空 select 永久阻塞
go
package main
func main() {
select {} // 永久阻塞,常用于让主 goroutine 不退出
}空 select 会永久阻塞,因为没有 case 能就绪。这在某些场景下有用(比如让程序保持运行),但要小心死锁。
2. select 中可以有空 case 体
go
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan struct{})
go func() {
time.Sleep(time.Second)
close(ch)
}()
select {
case <-ch:
// 只关心 ch 是否就绪,不需要做什么
}
fmt.Println("ch 就绪了")
}3. for + select 循环监听
go
package main
import (
"fmt"
"time"
)
func main() {
ticker := time.NewTicker(300 * time.Millisecond)
defer ticker.Stop()
done := make(chan struct{})
go func() {
time.Sleep(2 * time.Second)
close(done)
}()
loop:
for {
select {
case <-ticker.C:
fmt.Println("tick", time.Now().Format("15:04:05"))
case <-done:
fmt.Println("done")
break loop // 用 label 跳出 for,否则 break 只跳出 select
}
}
}注意:break 在 select 内只会跳出 select,不会跳出外层 for。要跳出 for 需要用标签(label)或 return。
十二、小结
本篇我们深入学习了 select 语句,它是 Go 并发编程的核心工具:
多路复用:
select能同时监听多个 channel 操作,哪个就绪执行哪个。随机选择:多个 case 同时就绪时,
select随机选一个(而非顺序),保证公平性,避免饥饿。这是与switch的关键区别。default 分支:实现非阻塞操作——所有 case 不就绪时执行 default,常用于非阻塞收发、轮询。
超时控制:
select + time.After是实现超时的标准模式。高频场景用time.NewTimer避免time.After的内存问题。ticker 和 timer:
time.NewTicker周期触发,time.NewTimer单次触发,都用完要Stop。非阻塞读写:
select + default实现「试一下不行就算了」的语义,还能实现优先级 channel。nil channel:对 nil channel 的操作永久阻塞,在 select 中等同于「禁用该 case」。利用这点可动态启用/禁用监听。
退出信号模式:goroutine 用
select监听 done channel,收到信号就退出。这是 context 包的基础。扇入模式:合并多个 channel 为一个,消费者只需监听一个,简化设计。
细节:空 select 永久阻塞;
break在 select 内只跳出 select 不跳出 for,需用 label;for + select是循环监听的标准写法。
下一篇我们将转向 sync 包,学习 WaitGroup、Once、Mutex 等传统同步原语,它们与 channel 互补,是处理共享状态的有力工具。