Appearance
高级并发: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/errgroup2. 基本用法
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,slow 和 veryslow 通过 select <-ctx.Done() 感知到取消并退出,不必白等 3 秒、5 秒。这就是 errgroup + context 的威力——首个错误触发全局取消。
二、errgroup vs WaitGroup
| 特性 | sync.WaitGroup | errgroup.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:123、article: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 轻量性 |
| 02 | goroutine 生命周期 | go 启动、WaitGroup、泄漏检测、数量控制 |
| 03 | Channel 基础 | 无缓冲/有缓冲、关闭、for range、方向限制 |
| 04 | select 多路复用 | 随机选择、default、超时、退出信号、扇入 |
| 05 | sync 包 | WaitGroup、Once、Mutex、RWMutex、Cond、Map、Pool |
| 06 | 并发模式 | WorkerPool、Pipeline、Fan-out/Fan-in、Generator、Future |
| 07 | context 包 | 取消、超时、传值、级联取消、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 |
| 等待一组 goroutine | WaitGroup / errgroup |
| 取消/超时传播 | context |
| 控制并发数 | semaphore / errgroup.SetLimit |
| 多 channel 监听 | select |
| 对象复用 | sync.Pool |
| 单例/一次性初始化 | sync.Once |
| 防缓存击穿 | singleflight |
| 配置热更新 | atomic.Value |
4. 进阶学习路径
- 深入源码:阅读
runtime/chan.go、runtime/select.go、sync/mutex.go,理解底层实现。 - 并发测试:学习
goleak、go test -race,写出能检测并发问题的测试。 - 性能调优:用 pprof 分析 goroutine、锁竞争,优化热点。
- 实践项目:
- 实现一个并发安全的 LRU 缓存。
- 写一个带超时和重试的 HTTP 客户端。
- 实现一个并发爬虫,限制并发数、支持取消。
- 用 Pipeline 模式处理流式数据。
- 关注新特性:Go 每个版本都在改进并发能力(如 Go 1.22 循环变量、Go 1.23 range over func),保持学习。
5. 推荐阅读
- Effective Go 的并发章节(官方)。
- Go 并发编程实战(书籍)。
- Go 官方博客关于 channel、context、memory model 的文章。
golang.org/x/sync的源码(很短很清晰)。
十三、小结
本篇是 Go 并发编程系列的收官,我们学习了:
errgroup:带错误处理的 goroutine 组。
Go启动、Wait等待并返回第一个错误、WithContext实现首个错误触发取消、SetLimit限制并发。比 WaitGroup 更强大,需要错误处理或取消传播时首选。semaphore:加权信号量,控制并发数。
Acquire获取令牌、Release释放、支持权重和超时。比 channel 做信号量更灵活,适合不同任务消耗不同资源的场景。singleflight:单飞模式,合并相同 key 的并发请求,防缓存击穿。
Do同步、DoChan异步。错误也共享,适合读操作合并。sync.Once 高级用法:带重试的初始化(自定义 done 标志)、一次性关闭(防止重复 Close)。
实战案例:并发抓取 URL(errgroup)、批量处理限流(semaphore)、缓存防击穿(singleflight)、综合处理器(errgroup + context + SetLimit)。
总结:Go 并发编程的核心是 goroutine + channel + select + sync + context + atomic + 扩展库。心法是「通信优于共享、正确性优先、避免泄漏、控制并发、优雅退出」。工具选择根据需求:传递数据用 channel、保护状态用锁、计数用 atomic、取消用 context、限流用 semaphore、防击穿用 singleflight。
并发编程是 Go 最精华的部分,也是区分初级和高级 Go 开发者的试金石。希望这个系列能帮你建立扎实的并发编程能力。接下来,多写、多测、多看源码,你会在实践中不断加深理解。祝你在 Go 并发的世界里游刃有余!