深入理解 Go 并发编程模式
深入理解 Go 并发编程模式
并发是 Go 语言最强大的特性之一。和很多语言把并发当作"附加功能"不同,Go 从设计之初就把并发作为核心能力。今天我来分享几个日常开发中最实用的并发模式。
Goroutine 和 Channel:基础回顾
Goroutine 是由 Go 运行时管理的轻量级线程,Channel 则是连接它们的管道。最简单的例子:
func main() {
ch := make(chan string)
go func() {
ch <- "来自 goroutine 的消息"
}()
msg := <-ch
fmt.Println(msg)
}这很直观,但实际项目中需要更精巧的模式。
模式一:扇出 / 扇入(Fan-Out / Fan-In)
当有一个可以并行化的 CPU 密集任务时,扇出将工作分配给多个 goroutine,扇入则负责收集结果。
func fanOut(input <-chan int, workers int) []<-chan int {
channels := make([]<-chan int, workers)
for i := 0; i < workers; i++ {
channels[i] = process(input)
}
return channels
}
func fanIn(channels ...<-chan int) <-chan int {
var wg sync.WaitGroup
merged := make(chan int)
for _, ch := range channels {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for val := range c {
merged <- val
}
}(ch)
}
go func() {
wg.Wait()
close(merged)
}()
return merged
}在做微服务开发时,我经常用这个模式来并行聚合多个数据源——比如同时拉取用户资料、订单记录和个性化推荐。
模式二:工作池(Worker Pool)
工作池用来限制并发数量。当你调用有连接数限制的外部 API 或数据库时,这个模式必不可少。
func workerPool(jobs <-chan Job, results chan<- Result, numWorkers int) {
var wg sync.WaitGroup
for i := 0; i < numWorkers; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for job := range jobs {
result := processJob(job)
results <- result
}
}(i)
}
go func() {
wg.Wait()
close(results)
}()
}这里的关键在于:jobs channel 天然充当了任务队列的角色。Worker 空闲了就自动从中取任务,实现了自动负载均衡。
模式三:Context 实现取消控制
context 包是 Go 处理取消、超时和请求级别参数传递的标准方式。任何耗时操作都应该接受 context 参数。
func fetchData(ctx context.Context, url string) ([]byte, error) {
req, err := http.NewRequestWithContext(ctx, "GET", url, nil)
if err != nil {
return nil, err
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
return io.ReadAll(resp.Body)
}
// 带超时的调用方式
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
data, err := fetchData(ctx, "https://api.example.com/data")
if err != nil {
log.Fatal(err)
}
fmt.Println(string(data))
}模式四:Select 多路复用
select 语句可以同时等待多个 channel 操作。配合 time.After 使用,非常适合实现超时控制和心跳检测。
func processWithTimeout(input <-chan Data) {
for {
select {
case data := <-input:
handle(data)
case <-time.After(30 * time.Second):
log.Println("30 秒未收到数据,执行健康检查...")
healthCheck()
}
}
}模式五:流水线(Pipeline)
流水线将多个处理阶段串联起来,每个阶段是一组运行相同函数的 goroutine。每个阶段通过 channel 接收上游数据、向下游发送处理结果。
func generate(nums ...int) <-chan int {
out := make(chan int)
go func() {
for _, n := range nums {
out <- n
}
close(out)
}()
return out
}
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
out <- n * n
}
close(out)
}()
return out
}
func main() {
ch := generate(2, 3, 4)
out := square(ch)
for result := range out {
fmt.Println(result) // 4, 9, 16
}
}常见踩坑点
Goroutine 泄漏 — 务必确保 goroutine 有退出路径。用 context 取消或 done channel 来通知退出。
竞态条件 — 开发时养成加 go run -race 的习惯,它能在运行时捕获大部分数据竞争问题。
Channel 死锁 — 如果所有 goroutine 都在等待 channel 操作,Go 会 panic 并提示"all goroutines are asleep"。可以用带缓冲的 channel 或者重新梳理流水线结构来解决。
总结
Go 的并发模型之所以强大,正是因为它足够简洁。Goroutine 开销极低,Channel 保证线程安全,select 提供灵活的控制流。掌握这五个模式,就能应对大多数并发编程场景。
核心原则:不要通过共享内存来通信,而要通过通信来共享内存。
祝编码愉快!