目录

Go Select 与并发模式

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

https://img.zhaojq.top/20260729163100010.png
并发模式

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("轮询结束")
}

defaultselect 变成非阻塞操作。如果所有 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 包使用可以更方便地管理超时和取消。