Go Select 与并发模式
select 语句是 Go 并发编程中的多路复用器,允许同时监听多个 channel 操作。配合 goroutine 和 channel,select 可以实现丰富的并发模式,如超时控制、扇出扇入、管道处理等。掌握这些模式是编写健壮并发程序的关键。
select 语句基础
select 类似于 switch,但专门用于处理 channel 操作。每个 case 对应一个 channel 的发送或接收操作。
package main
import (
"fmt"
"time"
)
func main() {
ch1 := make(chan string)
ch2 := make(chan string)
go func() {
time.Sleep(100 * time.Millisecond)
ch1 <- "来自 ch1 的消息"
}()
go func() {
time.Sleep(200 * time.Millisecond)
ch2 <- "来自 ch2 的消息"
}()
// 接收两次
for i := 0; i < 2; i++ {
select {
case msg := <-ch1:
fmt.Println(msg)
case msg := <-ch2:
fmt.Println(msg)
}
}
}select 会阻塞直到某个 case 的 channel 操作可以执行。如果多个 case 同时就绪,随机选择一个执行。
default 分支(非阻塞操作)
default 分支在没有任何 case 就绪时立即执行,实现非阻塞的 channel 操作。
package main
import "fmt"
func main() {
ch := make(chan int, 5)
// 非阻塞发送
select {
case ch <- 42:
fmt.Println("发送成功")
default:
fmt.Println("channel 已满,发送被跳过")
}
// 非阻塞接收
select {
case val := <-ch:
fmt.Println("接收到:", val)
default:
fmt.Println("channel 为空,接收被跳过")
}
// 轮询模式
ch2 := make(chan int, 3)
ch2 <- 1
ch2 <- 2
for {
select {
case val := <-ch2:
fmt.Printf("处理: %d\n", val)
default:
fmt.Println("没有更多数据,退出轮询")
goto done
}
}
done:
fmt.Println("轮询结束")
}default 让 select 变成非阻塞操作。如果所有 case 的 channel 都不可就绪,立即执行 default 分支。
超时控制
使用 time.After 配合 select 实现超时控制,这是 Go 中最常见的超时模式。
package main
import (
"fmt"
"time"
)
// 模拟一个可能超时的操作
func slowOperation(ch chan string) {
time.Sleep(2 * time.Second)
ch <- "操作完成"
}
func main() {
ch := make(chan string)
go slowOperation(ch)
// 设置 1 秒超时
select {
case result := <-ch:
fmt.Println("成功:", result)
case <-time.After(1 * time.Second):
fmt.Println("操作超时!")
}
// 带超时的函数
fmt.Println("\n--- 带超时的函数 ---")
result, err := fetchWithTimeout("http://example.com", 500*time.Millisecond)
if err != nil {
fmt.Println("错误:", err)
} else {
fmt.Println("结果:", result)
}
}
func fetchWithTimeout(url string, timeout time.Duration) (string, error) {
ch := make(chan string, 1) // 有缓冲,防止 goroutine 泄漏
go func() {
// 模拟网络请求
time.Sleep(1 * time.Second)
ch <- "响应数据"
}()
select {
case result := <-ch:
return result, nil
case <-time.After(timeout):
return "", fmt.Errorf("请求 %s 超时(%v)", url, timeout)
}
}time.After(d) 返回一个 channel,在持续时间 d 后自动发送当前时间。这是实现超时的惯用法。注意:使用有缓冲 channel 可以避免 goroutine 泄漏。
扇出模式(Fan-out)
扇出指将一个输入 channel 的数据分发给多个 worker 并行处理。
package main
import (
"fmt"
"sync"
"time"
)
// worker 处理任务并返回结果
func worker(id int, jobs <-chan int, results chan<- int, wg *sync.WaitGroup) {
defer wg.Done()
for job := range jobs {
fmt.Printf("Worker %d 开始处理任务 %d\n", id, job)
time.Sleep(100 * time.Millisecond) // 模拟处理时间
results <- job * 2
}
}
func main() {
jobs := make(chan int, 100)
results := make(chan int, 100)
var wg sync.WaitGroup
// 启动 3 个 worker(扇出)
for w := 1; w <= 3; w++ {
wg.Add(1)
go worker(w, jobs, results, &wg)
}
// 发送 9 个任务
for j := 1; j <= 9; j++ {
jobs <- j
}
close(jobs) // 关闭 jobs,worker 的 range 循环会结束
// 等待所有 worker 完成后关闭 results
go func() {
wg.Wait()
close(results)
}()
// 收集结果
for result := range results {
fmt.Printf("结果: %d\n", result)
}
fmt.Println("所有任务处理完毕")
}多个 worker 从同一个 channel 读取数据,自动实现负载均衡——哪个 worker 空闲就先拿到任务。
扇入模式(Fan-in)
扇入指将多个输入 channel 的数据合并到一个 channel。
package main
import (
"fmt"
"sync"
"time"
)
// 生产者:每个 goroutine 产生一组数据
func producer(id int, count int) <-chan int {
ch := make(chan int)
go func() {
for i := 0; i < count; i++ {
ch <- id*100 + i
time.Sleep(50 * time.Millisecond)
}
close(ch)
}()
return ch
}
// merge 将多个 channel 合并为一个(扇入)
func merge(channels ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
// 将每个 channel 的数据复制到 out
output := func(ch <-chan int) {
defer wg.Done()
for val := range ch {
out <- val
}
}
for _, ch := range channels {
wg.Add(1)
go output(ch)
}
// 所有输入 channel 关闭后,关闭输出 channel
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
// 创建 3 个生产者
ch1 := producer(1, 3)
ch2 := producer(2, 3)
ch3 := producer(3, 3)
// 扇入合并
merged := merge(ch1, ch2, ch3)
// 消费合并后的数据
for val := range merged {
fmt.Printf("收到: %d\n", val)
}
fmt.Println("所有生产者数据已合并处理完毕")
}扇入模式常用于聚合多个数据源的结果。merge 函数启动多个 goroutine 分别读取各个输入 channel,将数据统一写入一个输出 channel。
管道模式(Pipeline)
管道将数据处理分成多个阶段,每个阶段通过 channel 连接。
package main
import (
"fmt"
"math"
)
// 阶段 1:生成数字
func generate(nums ...int) <-chan int {
out := make(chan int)
go func() {
for _, n := range nums {
out <- n
}
close(out)
}()
return out
}
// 阶段 2:计算平方
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
out <- n * n
}
close(out)
}()
return out
}
// 阶段 3:过滤大于 10 的值
func filter(in <-chan int, threshold int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
if n > threshold {
out <- n
}
}
close(out)
}()
return out
}
// 阶段 4:计算平方根
func sqrt(in <-chan int) <-chan float64 {
out := make(chan float64)
go func() {
for n := range in {
out <- math.Sqrt(float64(n))
}
close(out)
}()
return out
}
func main() {
// 构建管道:生成 -> 平方 -> 过滤 -> 平方根
nums := generate(1, 2, 3, 4, 5, 6, 7)
squared := square(nums)
filtered := filter(squared, 10)
result := sqrt(filtered)
// 消费最终结果
fmt.Println("管道处理结果:")
for val := range result {
fmt.Printf(" %.2f\n", val)
}
}管道模式将复杂的数据处理拆分为多个独立的阶段,每个阶段职责单一,通过 channel 串联。这种模式易于组合、测试和维护。
done channel 模式
使用 done channel 通知 goroutine 退出,是 Go 中取消操作的标准模式。
package main
import (
"fmt"
"math/rand"
"time"
)
// 可取消的数据生成器
func generateNumbers(done <-chan struct{}) <-chan int {
ch := make(chan int)
go func() {
defer close(ch)
for {
select {
case <-done:
fmt.Println("生成器收到退出信号")
return
case ch <- rand.Intn(100):
// 发送成功,继续循环
}
}
}()
return ch
}
func main() {
done := make(chan struct{})
// 启动生成器
nums := generateNumbers(done)
// 只消费前 5 个数据
count := 0
for val := range nums {
fmt.Printf("收到: %d\n", val)
count++
if count >= 5 {
break
}
}
// 通知生成器退出
close(done)
time.Sleep(100 * time.Millisecond)
fmt.Println("程序结束")
}done channel 模式的核心思想:调用方通过关闭 done channel 来通知工作 goroutine 退出。工作 goroutine 在 select 中同时监听 done 和数据 channel,确保能及时响应取消信号。
综合示例:并发搜索引擎
下面综合使用多种模式实现一个简单的并发搜索。
package main
import (
"fmt"
"sync"
"time"
)
// SearchResult 搜索结果
type SearchResult struct {
Source string
Data string
}
// 模拟搜索不同数据源
func searchGoogle(query string, timeout time.Duration) SearchResult {
time.Sleep(150 * time.Millisecond)
return SearchResult{Source: "Google", Data: fmt.Sprintf("Google 结果: %s", query)}
}
func searchBing(query string, timeout time.Duration) SearchResult {
time.Sleep(200 * time.Millisecond)
return SearchResult{Source: "Bing", Data: fmt.Sprintf("Bing 结果: %s", query)}
}
func searchLocal(query string, timeout time.Duration) SearchResult {
time.Sleep(50 * time.Millisecond)
return SearchResult{Source: "Local", Data: fmt.Sprintf("本地结果: %s", query)}
}
// 并发搜索所有数据源,返回最先完成的结果
func concurrentSearch(query string) SearchResult {
ch := make(chan SearchResult, 3) // 有缓冲避免 goroutine 泄漏
go func() { ch <- searchGoogle(query, time.Second) }()
go func() { ch <- searchBing(query, time.Second) }()
go func() { ch <- searchLocal(query, time.Second) }()
// 返回最先到达的结果
return <-ch
}
// 并发搜索并收集所有结果
func searchAll(query string) []SearchResult {
ch := make(chan SearchResult, 3)
var wg sync.WaitGroup
sources := []func(string, time.Duration) SearchResult{
searchGoogle, searchBing, searchLocal,
}
for _, search := range sources {
wg.Add(1)
go func(fn func(string, time.Duration) SearchResult) {
defer wg.Done()
ch <- fn(query, time.Second)
}(search)
}
go func() {
wg.Wait()
close(ch)
}()
var results []SearchResult
for r := range ch {
results = append(results, r)
}
return results
}
func main() {
query := "Go 并发编程"
fmt.Println("=== 最快结果 ===")
fastest := concurrentSearch(query)
fmt.Printf("[%s] %s\n\n", fastest.Source, fastest.Data)
fmt.Println("=== 所有结果 ===")
all := searchAll(query)
for _, r := range all {
fmt.Printf("[%s] %s\n", r.Source, r.Data)
}
}总结
select 是 Go 并发编程的多路复用器,可以同时监听多个 channel 操作。default 分支实现非阻塞操作,time.After 实现超时控制。常见的并发模式包括:扇出(一个输入分给多个 worker)、扇入(多个输入合并到一个 channel)、管道(多阶段数据处理)和 done channel(取消通知)。这些模式可以灵活组合,构建高效的并发系统。在实际项目中,配合 context 包使用可以更方便地管理超时和取消。