本节把 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 限流、熔断与隔离 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。