《Go 语言编程实战》9.1 errgroup 与结构化并发

用 golang.org/x/sync/errgroup 把 TaskHub 的并发收进结构化边界:首个错误自动取消其余任务、SetLimit 限制并发数、TryGo 非阻塞提交,实测失败取消与并发上限,并对比裸 WaitGroup 在错误传播上的缺口。

本节把 TaskHub 推进到「并发可控」:一个详情页要并行拉任务、负责人、评论三份数据,本节用 errgroup 把它们的生命周期绑在一起——任一失败即取消其余,并限制并发数不让下游被打垮。
适用版本:Go 1.27(实测 go1.27.0),golang.org/x/sync。

9.1 errgroup 与结构化并发

TaskHub 的详情页要聚合多份数据:任务本体、负责人信息、评论列表。串行拉取是三个 RTT 相加,并行拉取只取最慢的那个。卷一讲过 goroutine 与 channel 的基础,本节讲工程上怎么把并发写得不会泄漏、不会失控——即结构化并发。

9.1.1 结构化并发:让子任务的生命周期受父任务约束

结构化并发的核心思想一句话:并发派生的任务,必须在父任务返回前全部结束。就像函数里的局部变量,作用域结束就该销毁。它带来三个保证:

  • 不会有孤儿 goroutine:父任务返回时,子任务一定已经退出(或已被取消)。
  • 错误能向上传播:子任务的失败能被父任务感知并处理。
  • 取消能向下传播:父任务取消时,子任务收到信号并退出。

Go 标准库的 WaitGroup 只给了「等待」这一个能力,缺了错误传播与取消。errgroup 补上了这两块。

9.1.2 裸 WaitGroup 的缺口

用 WaitGroup 并行拉三份数据,典型的写法是这样:

var wg sync.WaitGroup
var task Task
var owner User
var comments []Comment
wg.Add(3)
go func() { defer wg.Done(); task, _ = loadTask(ctx, id) }()
go func() { defer wg.Done(); owner, _ = loadOwner(ctx, id) }()
go func() { defer wg.Done(); comments, _ = loadComments(ctx, id) }()
wg.Wait()

问题一目了然:

  • 错误被丢掉了:每个 goroutine 里的 _ 把错误吞了,父任务不知道谁失败了。
  • 失败不会取消其余:如果 loadOwner 失败,loadTask、loadComments 依然跑到底,白白浪费资源。
  • 收集错误要自己写:想把三个错误都收上来,得自己加 mutex 或 channel,容易写错。

9.1.3 errgroup.WithContext:首个错误即取消

errgroup 的 WithContext 返回一个组和一个派生 context,ctx 会在任一任务返回错误或 Wait 返回时被取消:

g, ctx := errgroup.WithContext(context.Background())
g.Go(func() error { return loadTask(ctx, id) })
g.Go(func() error { return loadOwner(ctx, id) })
g.Go(func() error { return loadComments(ctx, id) })
if err := g.Wait(); err != nil {
	return err // 第一个非 nil 错误
}

实测「一个任务失败、其余任务被取消」的行为——5 个任务,第 2 个立刻返回错误,其余 4 个正在等待的会被 ctx 取消:

$ go run ./ch9/errgroup
task 5 被取消: context canceled
task 4 被取消: context canceled
task 1 被取消: context canceled
task 3 被取消: context canceled
Wait 返回: task 2 failed
被取消的任务数 = 4

4 个任务在 ctx.Done() 上醒来并退出,Wait 返回第一个错误 task 2 failed。这就是结构化并发:父任务一失败,子任务立刻收到取消信号,不留孤儿。Wait 只返回第一个错误(内部用 sync.Once 记录),后续错误被丢弃——这点很关键,见 9.1.5。

9.1.4 有界并发:SetLimit 与 TryGo

errgroup 默认不限制并发数——g.Go 每次都立刻起一个 goroutine。如果循环里对 10000 个任务各起一个 goroutine 去调下游,会把下游和本机都压垮。SetLimit(n) 限制同时活跃的 goroutine 上限:

g := new(errgroup.Group)
g.SetLimit(3)
for _, id := range ids {
	g.Go(func() error {
		return fetch(ctx, id) // 超过 3 个会阻塞在 Go 上,等有槽位再进
	})
}
_ = g.Wait()

实测并发上限设为 3 时,观测到的峰值并发正好是 3:

$ go run ./ch9/errgroup
并发上限设为 3,实测峰值并发 = 3

g.Go 在达到上限时阻塞,起到背压(backpressure)作用:生产速度自动降到消费速度。如果想要「满了就跳过、不阻塞」,用 TryGo——它在有空槽时返回 true、满了立刻返回 false:

if !g.TryGo(fn) {
	// 槽位已满,可以选择丢弃、排队或返回「服务繁忙」
}

实测 SetLimit(2) 下连续提交 6 个任务,只有 2 个被启动、4 个被立即拒绝:

$ go run ./ch9/trygo
SetLimit(2)+TryGo: 启动=2 立即拒绝=4

TryGo 适合「宁可拒绝也不要排队」的场景(如过载保护),Go 适合「排队也要处理完」的场景。选哪个取决于积压了会发生什么:积压可接受就 Go,积压会雪崩就 TryGo + 快速失败。

SetLimit 限制的是本组的 goroutine 数;如果要限制的是「对某个下游的全局并发」(多个 errgroup 共享一个预算),应该用 semaphore。两者分工不同:

sem := semaphore.NewWeighted(10) // 全局最多 10 个并发调用下游

g := new(errgroup.Group)
for _, id := range ids {
	g.Go(func() error {
		if err := sem.Acquire(ctx, 1); err != nil {
			return err
		}
		defer sem.Release(1)
		return fetch(ctx, id)
	})
}

SetLimit 管「本批任务并发多少」,semaphore 管「对某个资源的全局并发多少」。TaskHub 对同一个下游的调用统一走一个 semaphore,这样即使有多个并发入口,下游的并发上限也是可控的。

9.1.5 并行的收益:延迟从求和变成取最大

并行的意义是把「多个依赖的延迟之和」变成「最慢那个的延迟」。串行拉三份各 50ms 的数据要 150ms,并行只要 50ms,实测:

$ go run ./ch9/parallel
串行三份数据: 156ms
并行三份数据: 50ms

156ms 对 50ms,约 3 倍差距。依赖越多,收益越大——N 个等长依赖,串行是 N×T,并行是 T。但并行不是免费的:

  • 下游压力翻倍:原本错开的请求变成同时到达,下游瞬时并发变高,这就是要 SetLimit 的原因。
  • 调试更难:并发 bug(数据竞争、死锁)比串行难复现,go test -race 必须开(卷一 11.3)。
  • 错误路径更复杂:一个失败要取消其余,取消逻辑写错就泄漏。

所以判断标准是:依赖之间没有先后关系、且延迟可观(毫秒级以上)时才并行。三个内存里的 map 查询并行毫无意义,反而增加调度开销。50ms 以上的下游调用才值得并行。

9.1.6 只返回第一个错误 vs 收集全部

errgroup.Wait 只给第一个错误,对「快速失败」足够,但有时你需要知道所有任务各自的结果(比如批量校验,要一次性告诉用户所有问题)。两种做法:

需求方案
任一失败即整体失败errgroup 默认行为
收集所有错误每个 goroutine 把错误写进受 mutex 保护的切片,最后 errors.Join
部分失败可接受用 errgroup 但只在「致命错误」时返回非 nil

收集全部错误的写法:

var (
	mu   sync.Mutex
	errs []error
)
g, ctx := errgroup.WithContext(ctx)
for _, item := range items {
	g.Go(func() error {
		if err := validate(ctx, item); err != nil {
			mu.Lock()
			errs = append(errs, fmt.Errorf("item %v: %w", item, err))
			mu.Unlock()
		}
		return nil // 不返回错误,让其余任务继续
	})
}
_ = g.Wait()
if len(errs) > 0 {
	return errors.Join(errs...) // 卷一 6.2 的 errors.Join
}

注意这里 g.Go 里的函数故意返回 nil,否则第一个错误就会取消其余任务,收集不全。这是 errgroup 的灵活用法——它的取消能力可以按需关闭。

9.1.7 TaskHub 实战:并行聚合详情页

把上面的拼起来,TaskHub 的任务详情接口用 errgroup 并行拉三份数据,并限制对下游的并发:

func (s *Service) TaskDetail(ctx context.Context, id int64) (*Detail, error) {
	var (
		task     Task
		owner    User
		comments []Comment
	)
	g, ctx := errgroup.WithContext(ctx)
	g.Go(func() (err error) { task, err = s.repo.Task(ctx, id); return })
	g.Go(func() (err error) { owner, err = s.repo.Owner(ctx, id); return })
	g.Go(func() (err error) { comments, err = s.repo.Comments(ctx, id); return })
	if err := g.Wait(); err != nil {
		return nil, fmt.Errorf("task detail: %w", err)
	}
	return &Detail{Task: task, Owner: owner, Comments: comments}, nil
}

三个查询并行,任一失败其余立刻取消,错误带上 %w 包装后向上传播(卷一 6.2)。注意每个 g.Go 里的闭包各自把结果写进自己的变量,没有共享写入——这是 errgroup 安全使用的前提:goroutine 之间不通过共享变量通信,而是各自写独立变量,最后由 Wait 之后的主 goroutine 汇总。

9.1.8 常见坑

  • g.Go 里闭包捕获循环变量:Go 1.22 起循环变量每次迭代独立,但若捕获的是外部可变变量仍要小心,推荐参数传入或 x := x。
  • 忘记检查 g.Wait 的返回值:错误被静默吞掉,这是最常见的 bug。
  • 无上限并发:循环里对海量任务 g.Go,把下游打垮,务必 SetLimit。
  • SetLimit 在有任务运行时调用:文档明确禁止「有活跃 goroutine 时改上限」,要在启动前设好。
  • errgroup 里的 goroutine 泄漏:如果 fn 忽略了 ctx,取消信号进不去,任务照样跑完,失去了取消的意义。
  • 以为 Wait 返回所有错误:只返回第一个,要全集得自己收集。
  • 在 g.Go 里写共享变量:数据竞争,用独立变量 + Wait 后汇总,或用 channel。
  • WithContext 的 ctx 没传给下游:派生 ctx 必须一路传下去,否则取消传不到底层。

小结

  • 结构化并发 = 子任务生命周期受父任务约束:不留孤儿、错误上抛、取消下传。
  • errgroup.WithContext 在任一任务失败时取消其余,实测 5 个任务中 1 个失败、4 个被取消。
  • SetLimit 限制并发上限并阻塞背压,TryGo 满了立即拒绝,按「积压会不会雪崩」来选。
  • Wait 只返回第一个错误;要收集全部错误就让 fn 返回 nil 并自己聚合,最后 errors.Join。
  • goroutine 之间不共享写入,各自写独立变量,Wait 后汇总。

并发写对了,但并发本身会带来新风险:突发的流量可能压垮下游。下一节给 TaskHub 装上限流、熔断与隔离,让它在依赖变慢或过载时「优雅地退化」而不是「一起挂掉」。

阅读导航:上一节:8.3 定时任务与延迟队列 · 下一节:9.2 限流、熔断与隔离 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「golang」更多文章

  1. 《Go 语言编程实战》目录
  2. 《Go 语言编程实战》18.3 上线、观测与迭代
  3. 《Go 语言编程实战》18.2 故障演练