深入理解 Go 并发编程模式

并发是 Go 语言最强大的特性之一。和很多语言把并发当作"附加功能"不同,Go 从设计之初就把并发作为核心能力。今天我来分享几个日常开发中最实用的并发模式。

Goroutine 和 Channel:基础回顾

Goroutine 是由 Go 运行时管理的轻量级线程,Channel 则是连接它们的管道。最简单的例子:

go
func main() {
    ch := make(chan string)

    go func() {
        ch <- "来自 goroutine 的消息"
    }()

    msg := <-ch
    fmt.Println(msg)
}

这很直观,但实际项目中需要更精巧的模式。

模式一:扇出 / 扇入(Fan-Out / Fan-In)

当有一个可以并行化的 CPU 密集任务时,扇出将工作分配给多个 goroutine,扇入则负责收集结果。

go
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 或数据库时,这个模式必不可少。

go
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 参数。

go
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 使用,非常适合实现超时控制和心跳检测。

go
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 接收上游数据、向下游发送处理结果。

go
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 提供灵活的控制流。掌握这五个模式,就能应对大多数并发编程场景。

核心原则:不要通过共享内存来通信,而要通过通信来共享内存。

祝编码愉快!