# Go Select 与并发模式

`select` 语句是 Go 并发编程中的多路复用器，允许同时监听多个 channel 操作。配合 goroutine 和 channel，`select` 可以实现丰富的并发模式，如超时控制、扇出扇入、管道处理等。掌握这些模式是编写健壮并发程序的关键。

![并发模式](https://img.zhaojq.top/20260729163100010.png "并发模式")

## select 语句基础

`select` 类似于 `switch`，但专门用于处理 channel 操作。每个 `case` 对应一个 channel 的发送或接收操作。

```go
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 操作。

```go
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 中最常见的超时模式。

```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 并行处理。

```go
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。

```go
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 连接。

```go
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 中取消操作的标准模式。

```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，确保能及时响应取消信号。

## 综合示例：并发搜索引擎

下面综合使用多种模式实现一个简单的并发搜索。

```go
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` 包使用可以更方便地管理超时和取消。

