Skip to content

高级并发:errgroup 与扩展库

前面九章我们学习了标准库的并发工具。本篇是 Go 并发编程系列的收官之作,我们将学习 golang.org/x/sync 扩展库提供的三个利器:errgroup(带错误处理的 goroutine 组)、semaphore(信号量)、singleflight(单飞防击穿)。它们解决了标准库没覆盖的常见痛点,是生产级并发程序的标准配置。最后我们会总结整个并发编程的学习路径。

一、golang.org/x/sync/errgroup:带错误处理的 goroutine 组

回顾一下 sync.WaitGroup:它能等待一组 goroutine 完成,但不处理错误——goroutine 里出错了,主 goroutine 也无从得知。errgroup.Group 弥补了这个缺陷:它在 WaitGroup 之上增加了错误收集和取消传播能力。

1. 安装

bash
go get golang.org/x/sync/errgroup

2. 基本用法

go
package main

import (
	"context"
	"errors"
	"fmt"
	"time"

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

func fetch(url string) (string, error) {
	time.Sleep(time.Second) // 模拟网络请求
	if url == "bad" {
		return "", errors.New("获取 " + url + " 失败")
	}
	return "来自 " + url + " 的数据", nil
}

func main() {
	g := new(errgroup.Group)

	urls := []string{"api1", "api2", "bad", "api4"}
	results := make([]string, len(urls))

	for i, url := range urls {
		i, url := i, url // 捕获循环变量
		g.Go(func() error {
			data, err := fetch(url)
			if err != nil {
				return err // 返回错误
			}
			results[i] = data
			return nil
		})
	}

	// Wait 返回第一个非 nil 错误
	if err := g.Wait(); err != nil {
		fmt.Println("出错:", err)
	}
	fmt.Println("结果:", results)
}

3. errgroup 的核心特性

  • g.Go(fn func() error):启动一个 goroutine 执行 fn,fn 返回 error。
  • g.Wait():等待所有 goroutine 完成,返回第一个非 nil 错误。
  • 首个错误优先:一旦某个 goroutine 返回错误,后续的错误会被忽略(但 goroutine 仍会跑完)。
  • 自动 context 取消:配合 WithContext,任一 goroutine 出错会自动取消 context,通知其他 goroutine 停止。

4. 与 context 联合使用

go
package main

import (
	"context"
	"errors"
	"fmt"
	"time"

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

func main() {
	// WithContext 返回一个 group 和一个 context
	// 任一 goroutine 返回错误,context 会被自动取消
	g, ctx := errgroup.WithContext(context.Background())

	urls := []string{"fast", "slow", "bad", "veryslow"}

	for _, url := range urls {
		url := url
		g.Go(func() error {
			return worker(ctx, url)
		})
	}

	if err := g.Wait(); err != nil {
		fmt.Println("第一个错误:", err)
	} else {
		fmt.Println("全部成功")
	}
}

func worker(ctx context.Context, name string) error {
	delay := time.Second
	if name == "slow" {
		delay = 3 * time.Second
	} else if name == "veryslow" {
		delay = 5 * time.Second
	} else if name == "bad" {
		delay = 500 * time.Millisecond
	}

	select {
	case <-time.After(delay):
		if name == "bad" {
			return errors.New(name + " 失败")
		}
		fmt.Printf("%s 成功\n", name)
		return nil
	case <-ctx.Done():
		fmt.Printf("%s 被取消: %v\n", name, ctx.Err())
		return ctx.Err()
	}
}

关键点:bad 在 500ms 后返回错误,errgroup 自动取消 ctx,slowveryslow 通过 select <-ctx.Done() 感知到取消并退出,不必白等 3 秒、5 秒。这就是 errgroup + context 的威力——首个错误触发全局取消

二、errgroup vs WaitGroup

特性sync.WaitGrouperrgroup.Group
等待 goroutine
错误处理❌(需自己用 channel)✅(自动收集第一个错误)
取消传播✅(WithContext 自动取消)
限制并发数✅(SetLimit)
API 简洁度中(Add/Done/Wait)高(Go/Wait)

经验法则:需要错误处理或取消传播时用 errgroup,纯等待用 WaitGroup

三、golang.org/x/sync/semaphore:信号量控制并发数

semaphore.Weighted 是一个加权信号量,用于控制同时访问某资源的 goroutine 数量。比用带缓冲 channel 做信号量更灵活(支持不同权重)。

1. 基本用法

go
package main

import (
	"context"
	"fmt"
	"sync"
	"time"

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

func main() {
	// 创建容量为 3 的信号量
	sem := semaphore.NewWeighted(3)
	var wg sync.WaitGroup

	for i := 1; i <= 10; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			worker(sem, id)
		}(i)
	}

	wg.Wait()
	fmt.Println("全部完成")
}

func worker(sem *semaphore.Weighted, id int) {
	ctx := context.Background()
	// 获取一个令牌(权重 1)
	if err := sem.Acquire(ctx, 1); err != nil {
		fmt.Printf("worker %d 获取信号量失败: %v\n", id, err)
		return
	}
	defer sem.Release(1)

	fmt.Printf("worker %d 开始 (当前并发可继续)\n", id)
	time.Sleep(time.Second) // 模拟工作
	fmt.Printf("worker %d 完成\n", id)
}

任意时刻最多只有 3 个 goroutine 在执行 worker 主体,其余在 Acquire 处阻塞。

2. 带超时的 Acquire

go
package main

import (
	"context"
	"fmt"
	"time"

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

func main() {
	sem := semaphore.NewWeighted(2)

	// 先占满
	ctx := context.Background()
	sem.Acquire(ctx, 2)

	// 再尝试获取,带超时
	timeoutCtx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
	defer cancel()

	if err := sem.Acquire(timeoutCtx, 1); err != nil {
		fmt.Println("获取超时:", err) // context deadline exceeded
	}

	sem.Release(2)
	fmt.Println("释放完毕")
}

3. 加权信号量

go
package main

import (
	"context"
	"fmt"
	"sync"

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

func main() {
	// 总容量 10
	sem := semaphore.NewWeighted(10)
	var wg sync.WaitGroup

	// 大任务占 5 个令牌
	wg.Add(1)
	go func() {
		defer wg.Done()
		sem.Acquire(context.Background(), 5)
		defer sem.Release(5)
		fmt.Println("大任务执行(占5)")
	}()

	// 小任务各占 2 个令牌
	for i := 0; i < 3; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			sem.Acquire(context.Background(), 2)
			defer sem.Release(2)
			fmt.Printf("小任务 %d 执行(占2)\n", id)
		}(i)
	}

	wg.Wait()
	fmt.Println("全部完成")
}

加权信号量适合「不同任务消耗不同资源」的场景,比如大文件下载占 5 个连接、小文件占 2 个。

四、golang.org/x/sync/singleflight:单飞模式

singleflight.Group 实现「单飞」:对同一个 key 的多个并发请求,只真正执行一次,所有请求共享这一次的结果。这是防缓存击穿的经典武器。

1. 缓存击穿问题

缓存击穿:某个热点 key 过期的瞬间,大量请求同时打到数据库。如果用 singleflight,这些请求会被合并成一次实际查询,结果共享给所有请求。

2. 基本用法

go
package main

import (
	"fmt"
	"sync"
	"time"

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

var group singleflight.Group

// 模拟昂贵的数据库查询
func expensiveQuery(key string) (string, error) {
	time.Sleep(time.Second) // 模拟耗时
	return "数据-" + key, nil
}

func query(key string) (string, error) {
	// Do 会合并相同 key 的调用
	// 第一个调用执行 fn,其他调用阻塞等待结果
	v, err, shared := group.Do(key, func() (any, error) {
		fmt.Printf(">>> 实际查询 %s\n", key)
		return expensiveQuery(key)
	})
	if err != nil {
		return "", err
	}
	fmt.Printf("得到结果: %v (共享: %v)\n", v, shared)
	return v.(string), nil
}

func main() {
	var wg sync.WaitGroup

	// 10 个 goroutine 同时查同一个 key
	for i := 0; i < 10; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			result, _ := query("hotkey")
			fmt.Printf("goroutine %d: %s\n", id, result)
		}(i)
	}

	wg.Wait()
}

输出中「实际查询」只出现一次,但 10 个 goroutine 都拿到了结果。shared 为 true 表示结果被多个调用者共享。

3. singleflight 的注意事项

  • 只阻塞,不取消:等待中的调用者会一直等到第一个调用完成,不会因为某个调用者取消而中断。
  • 错误也共享:如果第一个调用返回错误,所有等待者都收到同样的错误。
  • 适合读操作:singleflight 用于合并「相同请求」,适合查询,不适合写操作(写需要每个都执行)。
  • key 选择:key 要能唯一标识一个请求,比如 user:123article:456

4. DoChan:异步版本

go
package main

import (
	"fmt"
	"sync"
	"time"

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

var group singleflight.Group

func main() {
	var wg sync.WaitGroup

	for i := 0; i < 5; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			// DoChan 返回一个 channel,结果就绪时发送
			ch := group.DoChan("key", func() (any, error) {
				time.Sleep(time.Second)
				return "结果", nil
			})
			result := <-ch
			fmt.Printf("goroutine %d: %v (shared: %v)\n", id, result.Val, result.Shared)
		}(i)
	}

	wg.Wait()
}

DoChan 适合需要 select 监听其他信号(如超时)的场景。

五、errgroup + context 联合使用详解

errgroup 的 WithContext 内部创建了一个 context,任一 goroutine 返回错误会自动取消它。我们再看一个更完整的实战。

1. 并行查询,首个失败即取消

go
package main

import (
	"context"
	"errors"
	"fmt"
	"math/rand"
	"time"

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

type Result struct {
	Name string
	Data string
}

func callService(ctx context.Context, name string, delay time.Duration, fail bool) (Result, error) {
	select {
	case <-time.After(delay):
		if fail {
			return Result{}, errors.New(name + " 失败")
		}
		return Result{Name: name, Data: name + " 数据"}, nil
	case <-ctx.Done():
		return Result{}, fmt.Errorf("%s 被取消", name)
	}
}

func main() {
	rand.Seed(time.Now().UnixNano())

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

	services := []struct {
		name  string
		delay time.Duration
		fail  bool
	}{
		{"用户服务", 300 * time.Millisecond, false},
		{"订单服务", 500 * time.Millisecond, false},
		{"支付服务", 200 * time.Millisecond, true}, // 这个会失败
		{"库存服务", 800 * time.Millisecond, false},
	}

	results := make([]Result, len(services))
	for i, svc := range services {
		i, svc := i, svc
		g.Go(func() error {
			r, err := callService(ctx, svc.name, svc.delay, svc.fail)
			if err != nil {
				return err
			}
			results[i] = r
			return nil
		})
	}

	if err := g.Wait(); err != nil {
		fmt.Println("失败:", err)
	}
	for _, r := range results {
		if r.Name != "" {
			fmt.Printf("成功: %+v\n", r)
		}
	}
}

支付服务在 200ms 后失败,errgroup 自动取消 ctx,库存服务(800ms)通过 select <-ctx.Done() 提前感知并退出,不必白等。

六、errgroup 的 SetLimit:限制并发数

Go 1.18 起,errgroup 支持 SetLimit(n) 限制同时运行的 goroutine 数量。

1. 用法

go
package main

import (
	"fmt"
	"sync"
	"time"

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

func main() {
	g := new(errgroup.Group)
	g.SetLimit(3) // 最多 3 个 goroutine 同时运行

	var mu sync.Mutex
	active := 0
	maxActive := 0

	for i := 1; i <= 10; i++ {
		i := i
		g.Go(func() error {
			mu.Lock()
			active++
			if active > maxActive {
				maxActive = active
			}
			mu.Unlock()

			fmt.Printf("任务 %d 开始\n", i)
			time.Sleep(300 * time.Millisecond)
			fmt.Printf("任务 %d 结束\n", i)

			mu.Lock()
			active--
			mu.Unlock()
			return nil
		})
	}

	g.Wait()
	fmt.Printf("最大并发数: %d\n", maxActive) // 3
}

g.Go 在达到限制时会阻塞,直到有 goroutine 完成腾出位置。这让 errgroup 同时具备了「错误处理 + 取消传播 + 并发限制」三大能力,非常强大。

2. TryGo:非阻塞启动

go
package main

import (
	"fmt"

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

func main() {
	g := new(errgroup.Group)
	g.SetLimit(2)

	// TryGo 非阻塞,达到限制返回 false
	for i := 0; i < 5; i++ {
		if ok := g.TryGo(func() error {
			return nil
		}); !ok {
			fmt.Printf("第 %d 个被拒绝(达到限制)\n", i+1)
		}
	}
	g.Wait()
	fmt.Println("完成")
}

七、sync.Once 的高级用法

sync.Once 除了单例,还有一些高级用法。

1. 带错误的重试初始化

Once 只执行一次,即使失败也不重试。如果需要失败可重试,可以自己实现:

go
package main

import (
	"errors"
	"fmt"
	"sync"
)

type OnceWithRetry struct {
	mu  sync.Mutex
	done bool
}

func (o *OnceWithRetry) Do(f func() error) error {
	o.mu.Lock()
	defer o.mu.Unlock()
	if o.done {
		return nil
	}
	if err := f(); err != nil {
		return err // 失败不标记 done,可重试
	}
	o.done = true
	return nil
}

func main() {
	var o OnceWithRetry
	attempts := 0

	for {
		attempts++
		err := o.Do(func() error {
			if attempts < 3 {
				return errors.New("还没准备好")
			}
			fmt.Println("初始化成功")
			return nil
		})
		if err == nil {
			break
		}
		fmt.Printf("第 %d 次尝试失败: %v\n", attempts, err)
	}
}

2. 用 Once 实现一次性关闭

go
package main

import (
	"fmt"
	"sync"
)

type Resource struct {
	closeOnce sync.Once
	closed    bool
}

func (r *Resource) Close() {
	r.closeOnce.Do(func() {
		r.closed = true
		fmt.Println("资源已关闭")
	})
}

func main() {
	r := &Resource{}
	r.Close()
	r.Close() // 不会重复关闭
	r.Close()
	fmt.Println("closed:", r.closed)
}

八、实战案例:并发抓取多个 URL

用 errgroup 并发抓取多个 URL,任一失败不影响其他,最后汇总结果。

go
package main

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

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

type FetchResult struct {
	URL    string
	Length int
	Err    error
}

func fetchURL(ctx context.Context, url string) (int, error) {
	req, err := http.NewRequestWithContext(ctx, "GET", url, nil)
	if err != nil {
		return 0, err
	}
	resp, err := http.DefaultClient.Do(req)
	if err != nil {
		return 0, err
	}
	defer resp.Body.Close()
	if resp.StatusCode != http.StatusOK {
		return 0, fmt.Errorf("状态码 %d", resp.StatusCode)
	}
	data, err := io.ReadAll(resp.Body)
	if err != nil {
		return 0, err
	}
	return len(data), nil
}

func main() {
	ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
	defer cancel()

	g, ctx := errgroup.WithContext(ctx)

	urls := []string{
		"https://golang.org",
		"https://httpbin.org/get",
		"https://httpbin.org/delay/1",
		// "https://invalid-url-xxx.com", // 取消注释可测试失败
	}

	results := make([]FetchResult, len(urls))
	var mu sync.Mutex

	for i, url := range urls {
		i, url := i, url
		g.Go(func() error {
			length, err := fetchURL(ctx, url)
			mu.Lock()
			results[i] = FetchResult{URL: url, Length: length, Err: err}
			mu.Unlock()
			// 这里不返回 err,让所有 URL 都尝试
			// 如果想首个失败即取消,return err 即可
			return nil
		})
	}

	if err := g.Wait(); err != nil {
		fmt.Println("错误:", err)
	}

	for _, r := range results {
		if r.Err != nil {
			fmt.Printf("[失败] %s: %v\n", r.URL, r.Err)
		} else {
			fmt.Printf("[成功] %s: %d 字节\n", r.URL, r.Length)
		}
	}
}

注意这里 g.Go 里始终返回 nil,是为了让所有 URL 都执行完(不因一个失败而取消其他)。如果需要「任一失败即停止」,把 return nil 改成 return err 即可。

九、实战案例:限制并发数的批量处理

用 semaphore 控制并发,批量处理大量任务而不压垮系统。

go
package main

import (
	"context"
	"fmt"
	"sync"
	"sync/atomic"
	"time"

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

func processItem(id int) error {
	time.Sleep(200 * time.Millisecond) // 模拟处理
	return nil
}

func main() {
	const totalItems = 20
	const maxConcurrent = 5

	sem := semaphore.NewWeighted(maxConcurrent)
	var wg sync.WaitGroup
	var success int64
	var failed int64

	start := time.Now()
	for i := 1; i <= totalItems; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()

			// 获取信号量
			if err := sem.Acquire(context.Background(), 1); err != nil {
				atomic.AddInt64(&failed, 1)
				return
			}
			defer sem.Release(1)

			if err := processItem(id); err != nil {
				atomic.AddInt64(&failed, 1)
				fmt.Printf("处理 %d 失败: %v\n", id, err)
				return
			}
			atomic.AddInt64(&success, 1)
			fmt.Printf("处理 %d 完成\n", id)
		}(i)
	}

	wg.Wait()
	elapsed := time.Since(start)
	fmt.Printf("\n成功 %d, 失败 %d, 总耗时 %v\n",
		success, failed, elapsed)
	// 20 个任务,每次最多 5 个并发,每个 200ms
	// 理论最少耗时: 4 批 * 200ms = 800ms
}

这个模式适合「批量调 API、批量下载、批量入库」等场景,控制并发避免压垮下游。

十、实战案例:缓存防击穿

用 singleflight 实现防击穿的缓存。

go
package main

import (
	"context"
	"fmt"
	"sync"
	"time"

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

type Cache struct {
	mu       sync.RWMutex
	data     map[string]string
	expireAt map[string]time.Time
	group    singleflight.Group
}

func NewCache() *Cache {
	return &Cache{
		data:     make(map[string]string),
		expireAt: make(map[string]time.Time),
	}
}

// 模拟从数据库加载
func loadFromDB(ctx context.Context, key string) (string, error) {
	time.Sleep(time.Second) // 模拟慢查询
	return "DB:" + key, nil
}

func (c *Cache) Get(ctx context.Context, key string) (string, error) {
	// 先查缓存
	c.mu.RLock()
	if val, ok := c.data[key]; ok {
		if time.Now().Before(c.expireAt[key]) {
			c.mu.RUnlock()
			return val, nil
		}
	}
	c.mu.RUnlock()

	// 缓存未命中,用 singleflight 合并请求
	v, err, _ := c.group.Do(key, func() (any, error) {
		fmt.Printf(">>> 实际加载 %s\n", key)
		val, err := loadFromDB(ctx, key)
		if err != nil {
			return nil, err
		}
		// 写入缓存
		c.mu.Lock()
		c.data[key] = val
		c.expireAt[key] = time.Now().Add(5 * time.Minute)
		c.mu.Unlock()
		return val, nil
	})
	if err != nil {
		return "", err
	}
	return v.(string), nil
}

func main() {
	cache := NewCache()
	ctx := context.Background()

	// 模拟缓存击穿:10 个并发请求同一个 key
	var wg sync.WaitGroup
	for i := 0; i < 10; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			val, _ := cache.Get(ctx, "user:123")
			fmt.Printf("goroutine %d: %s\n", id, val)
		}(i)
	}
	wg.Wait()

	fmt.Println("\n第二次查询(缓存命中)")
	val, _ := cache.Get(ctx, "user:123")
	fmt.Println("结果:", val)
}

第一次查询时,10 个 goroutine 同时请求,但 singleflight 只让一个真正查数据库,其他共享结果。第二次查询直接命中缓存,不再查库。这就是防击穿的核心思想。

十一、综合实战:errgroup + semaphore + context

把几个工具组合起来,实现一个生产级的并发任务处理器:支持并发限制、错误处理、取消传播、超时。

go
package main

import (
	"context"
	"errors"
	"fmt"
	"math/rand"
	"sync"
	"sync/atomic"
	"time"

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

type Task struct {
	ID    int
	Heavy bool
}

type Result struct {
	TaskID int
	Output string
	Err    error
}

func processTask(ctx context.Context, task Task) (string, error) {
	delay := 500 * time.Millisecond
	if task.Heavy {
		delay = 2 * time.Second
	}
	select {
	case <-time.After(delay):
		if rand.Intn(10) < 2 { // 20% 失败率
			return "", fmt.Errorf("任务 %d 失败", task.ID)
		}
		return fmt.Sprintf("任务 %d 完成", task.ID), nil
	case <-ctx.Done():
		return "", ctx.Err()
	}
}

func runTasks(ctx context.Context, tasks []Task, concurrency int) ([]Result, error) {
	g, ctx := errgroup.WithContext(ctx)
	g.SetLimit(concurrency) // errgroup 自带并发限制

	results := make([]Result, len(tasks))
	for i, task := range tasks {
		i, task := i, task
		g.Go(func() error {
			output, err := processTask(ctx, task)
			results[i] = Result{TaskID: task.ID, Output: output, Err: err}
			return err // 返回错误会触发 errgroup 取消
		})
	}

	if err := g.Wait(); err != nil {
		// 返回部分结果和第一个错误
		return results, err
	}
	return results, nil
}

func main() {
	rand.Seed(time.Now().UnixNano())

	// 构造任务
	tasks := make([]Task, 15)
	for i := range tasks {
		tasks[i] = Task{ID: i + 1, Heavy: i%5 == 0}
	}

	// 总超时 5 秒,并发 4
	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel()

	start := time.Now()
	results, err := runTasks(ctx, tasks, 4)
	elapsed := time.Since(start)

	var success, failed, canceled int64
	for _, r := range results {
		if r.Err == nil {
			success++
		} else if errors.Is(r.Err, context.Canceled) || errors.Is(r.Err, context.DeadlineExceeded) {
			canceled++
		} else {
			failed++
		}
	}

	fmt.Printf("耗时: %v\n", elapsed)
	fmt.Printf("成功: %d, 失败: %d, 取消: %d\n", success, failed, canceled)
	if err != nil {
		fmt.Println("首个错误:", err)
	}
}

这个例子集成了:

  • errgroup:错误收集、Wait 等待、并发限制(SetLimit)。
  • context:超时控制、取消传播。
  • select:每个任务监听 ctx.Done(),及时退出。

这是真实业务中「批量并发处理 + 容错 + 超时」的标准骨架。

十二、并发编程总结与学习路径

恭喜你完成了 Go 并发编程系列的全部内容!让我们做一个系统总结。

1. 知识体系回顾

章节主题核心要点
01并发基础并发 vs 并行、GMP 模型、goroutine 轻量性
02goroutine 生命周期go 启动、WaitGroup、泄漏检测、数量控制
03Channel 基础无缓冲/有缓冲、关闭、for range、方向限制
04select 多路复用随机选择、default、超时、退出信号、扇入
05sync 包WaitGroup、Once、Mutex、RWMutex、Cond、Map、Pool
06并发模式WorkerPool、Pipeline、Fan-out/Fan-in、Generator、Future
07context 包取消、超时、传值、级联取消、HTTP/DB 应用
08数据竞争race 检测、atomic、CAS、happens-before
09陷阱与最佳实践泄漏、关闭姿势、Mutex 复制、循环变量、优雅退出
10扩展库errgroup、semaphore、singleflight

2. 核心心法

  • 哲学:不要通过共享内存通信,而要通过通信共享内存。优先 channel,但该用锁时用锁。
  • 正确性优先:先保证正确,再优化性能。用 -race 检测竞态。
  • 避免泄漏:每个 goroutine 都要有明确的退出路径,用 context 传播取消。
  • 控制并发:不要无限创建 goroutine,用 Worker Pool 或信号量限制。
  • 优雅退出:监听信号、传播取消、Wait 等待、超时兜底。

3. 工具选择速查

需求推荐工具
goroutine 间传递数据channel
保护共享状态Mutex / RWMutex
简单计数器atomic
等待一组 goroutineWaitGroup / errgroup
取消/超时传播context
控制并发数semaphore / errgroup.SetLimit
多 channel 监听select
对象复用sync.Pool
单例/一次性初始化sync.Once
防缓存击穿singleflight
配置热更新atomic.Value

4. 进阶学习路径

  1. 深入源码:阅读 runtime/chan.goruntime/select.gosync/mutex.go,理解底层实现。
  2. 并发测试:学习 goleakgo test -race,写出能检测并发问题的测试。
  3. 性能调优:用 pprof 分析 goroutine、锁竞争,优化热点。
  4. 实践项目
    • 实现一个并发安全的 LRU 缓存。
    • 写一个带超时和重试的 HTTP 客户端。
    • 实现一个并发爬虫,限制并发数、支持取消。
    • 用 Pipeline 模式处理流式数据。
  5. 关注新特性:Go 每个版本都在改进并发能力(如 Go 1.22 循环变量、Go 1.23 range over func),保持学习。

5. 推荐阅读

  • Effective Go 的并发章节(官方)。
  • Go 并发编程实战(书籍)。
  • Go 官方博客关于 channel、context、memory model 的文章。
  • golang.org/x/sync 的源码(很短很清晰)。

十三、小结

本篇是 Go 并发编程系列的收官,我们学习了:

  1. errgroup:带错误处理的 goroutine 组。Go 启动、Wait 等待并返回第一个错误、WithContext 实现首个错误触发取消、SetLimit 限制并发。比 WaitGroup 更强大,需要错误处理或取消传播时首选。

  2. semaphore:加权信号量,控制并发数。Acquire 获取令牌、Release 释放、支持权重和超时。比 channel 做信号量更灵活,适合不同任务消耗不同资源的场景。

  3. singleflight:单飞模式,合并相同 key 的并发请求,防缓存击穿。Do 同步、DoChan 异步。错误也共享,适合读操作合并。

  4. sync.Once 高级用法:带重试的初始化(自定义 done 标志)、一次性关闭(防止重复 Close)。

  5. 实战案例:并发抓取 URL(errgroup)、批量处理限流(semaphore)、缓存防击穿(singleflight)、综合处理器(errgroup + context + SetLimit)。

  6. 总结:Go 并发编程的核心是 goroutine + channel + select + sync + context + atomic + 扩展库。心法是「通信优于共享、正确性优先、避免泄漏、控制并发、优雅退出」。工具选择根据需求:传递数据用 channel、保护状态用锁、计数用 atomic、取消用 context、限流用 semaphore、防击穿用 singleflight。

并发编程是 Go 最精华的部分,也是区分初级和高级 Go 开发者的试金石。希望这个系列能帮你建立扎实的并发编程能力。接下来,多写、多测、多看源码,你会在实践中不断加深理解。祝你在 Go 并发的世界里游刃有余!