并发模式:工作池与 errgroup

原语(goroutine、channel、锁)解决"怎么通信与同步",模式解决"如何组织"才能既高效又可控。本节介绍两个最常用的生产模式:worker pool 与 errgroup。

Worker Pool:固定大小工作池

无限制地 go func() 会在任务量大时耗尽内存或打垮下游。工作池模式预先启动固定数量的 worker,从任务通道持续取任务,实现并发度上限背压

package main

import (
    "fmt"
    "sync"
    "time"
)

type Job struct {
    ID   int
    Data string
}

type Result struct {
    JobID int
    Out   string
    Err   error
}

func worker(id int, jobs <-chan Job, results chan<- Result, wg *sync.WaitGroup) {
    defer wg.Done()
    for job := range jobs { // 通道关闭后循环自然结束
        // 模拟耗时处理
        time.Sleep(100 * time.Millisecond)
        results <- Result{JobID: job.ID, Out: "processed:" + job.Data}
    }
}

func main() {
    const numWorkers = 3
    const numJobs = 8

    jobs := make(chan Job, numJobs)
    results := make(chan Result, numJobs)

    // 启动固定数量的 worker
    var wg sync.WaitGroup
    for i := 1; i <= numWorkers; i++ {
        wg.Add(1)
        go worker(i, jobs, results, &wg)
    }

    // 投递任务
    for j := 1; j <= numJobs; j++ {
        jobs <- Job{ID: j, Data: fmt.Sprintf("task-%d", j)}
    }
    close(jobs) // 投递完毕,通知 worker 退出

    // 收集结果:先在后台等 worker 全部退出,再关闭 results
    go func() {
        wg.Wait()
        close(results)
    }()

    for r := range results {
        fmt.Println("worker 完成", r.JobID, r.Out)
    }
}

要点:

  • jobs 通道由投递方关闭,worker 通过 for range 感知退出。
  • results 由独立的收尾 goroutine 在 wg.Wait() 后关闭,主流程用 for range 安全收集。
  • 并发度 = worker 数量;任务缓冲用于平滑投递高峰。
Tip

worker 数量的经验值:CPU 密集型任务取 runtime.NumCPU();I/O 密集型(HTTP 调用、数据库访问)取决于下游承受能力,通常配额 10–100 并通过压测确定。

errgroup:带错误传播与限流的并发

批量并发任务中,任一子任务失败往往应当整体取消。golang.org/x/sync/errgroup 将 WaitGroup、错误收集与 context 取消打包:

go get golang.org/x/sync
package main

import (
    "context"
    "fmt"
    "net/http"
    "time"

    "golang.org/x/sync/errgroup"
)

func fetchStatus(ctx context.Context, client *http.Client, url string) error {
    req, _ := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
    resp, err := client.Do(req)
    if err != nil {
        return err
    }
    defer resp.Body.Close()
    if resp.StatusCode >= 400 {
        return fmt.Errorf("%s 返回 %d", url, resp.StatusCode)
    }
    return nil
}

func main() {
    client := &http.Client{Timeout: 5 * time.Second}
    urls := []string{
        "https://httpbin.org/status/200",
        "https://httpbin.org/status/500", // 会失败
        "https://httpbin.org/status/200",
    }

    g, ctx := errgroup.WithContext(context.Background())

    for _, u := range urls {
        g.Go(func() error { // 并发执行
            return fetchStatus(ctx, client, u)
        })
    }

    // Wait 阻塞至全部结束,返回第一个非 nil 错误
    if err := g.Wait(); err != nil {
        fmt.Println("批量任务失败:", err)
        return
    }
    fmt.Println("全部成功")
}

errgroup 的三项核心能力:

能力API说明
等待全部完成g.Wait()语义同 WaitGroup,但返回首个错误
错误时整体取消errgroup.WithContext任一任务失败即取消派生 ctx,其余任务可感知退出
限制并发数g.SetLimit(n)超过 n 时 Go 阻塞,等效于内置工作池

SetLimit:有界并发的批量处理

g, ctx := errgroup.WithContext(context.Background())
g.SetLimit(5) // 最多 5 个并发

for _, item := range items {
    item := item
    g.Go(func() error {
        return process(ctx, item)
    })
}
if err := g.Wait(); err != nil {
    // 处理失败
}

Go 1.22+ 中循环变量已按迭代隔离,item := item 的复制写法不再必要,保留仅为兼容旧版本。

模式选型

场景推荐
需要区分任务与结果、结果带元数据worker pool(channel 版)
批量子任务、任一失败即取消errgroup.WithContext
批量子任务、仅限并发数errgroup + SetLimit
任务失败需各自记录而非中断worker pool + Result.Err 字段

:::warning 两条共性陷阱:

  1. 不要在循环里无限 go:任务量不可控时一律使用池或 SetLimit
  2. 通道关闭权与发送方绑定:无论哪种模式,只由发送方关闭通道,避免向已关闭通道发送导致 panic。 :::

小结

  • worker pool 以固定 goroutine 数消费任务通道,提供并发上限与背压。
  • errgroup 合并等待、错误传播与 context 取消,SetLimit 提供内置限流。
  • 两者可组合:用 errgroup 管理 worker 生命周期是常见的工程化封装。