Go 基础体系 · 第 86/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。
Go 并发工具:errgroup、singleflight、semaphore 与 ants
本文以 Go 1.26.4 为基准,示例固定使用 golang.org/x/sync v0.17.0 与 github.com/panjf2000/ants/v2 v2.11.3。errgroup 协调一组返回错误的任务,singleflight 合并同 key 的同时调用,semaphore 限制加权在途量,ants 复用有界 worker。它们解决不同问题,不能统称为“协程池”后互换。
并发工具不能替业务决定任务是否幂等、失败是否取消、队列满时拒绝还是等待。本文以 API 状态、context、资源所有权和故障诊断为主线。
1. 先用四个问题区分工具
第一个问题是要“等待一组任务”还是“长期接纳任务”。前者用 errgroup/WaitGroup,后者才可能需要池。第二个问题是限制同时执行量还是合并相同工作:semaphore 控并发,singleflight 控重复。第三个问题是任务等价性:只有 key 和结果语义完全相同才能共享。第四个问题是资源瓶颈:数据库连接、内存、CPU 与第三方配额的单位可能不同。
一批任务 + 首错取消 -> errgroup
同 key 同时回源只做一次 -> singleflight
每项成本不同的在途限制 -> weighted semaphore
极大量短任务且已证实调度成本 -> ants
普通 go f() 加有界 channel 或 WaitGroup 经常已经足够。先写清接纳、取消、等待和关闭协议,再选能减少代码而不隐藏语义的工具。
2. errgroup 的生命周期和首错语义
errgroup.WithContext(parent) 返回 group 和派生 context。任一 Go 函数首次返回非 nil error 时,派生 context 被取消;Wait 返回第一个非 nil error,并在所有已启动函数返回后结束。即使全部成功,Wait 返回时派生 context 也会取消,因此它只应服务该任务组。
func loadDashboard(ctx context.Context, ids []int64) error {
group, groupCtx := errgroup.WithContext(ctx)
group.SetLimit(8)
for _, id := range ids {
id := id
group.Go(func() error {
if err := loadWidget(groupCtx, id); err != nil {
return fmt.Errorf("load widget %d: %w", id, err)
}
return nil
})
}
if err := group.Wait(); err != nil {
return fmt.Errorf("load dashboard: %w", err)
}
return nil
}
任务函数必须把 groupCtx 传给 HTTP、SQL 和阻塞 select,并在取消时返回。取消只是信号,不会强制杀死不响应 context 的库调用。Wait 是所有权屏障:它保证任务函数已返回,之后调用者才能释放它们使用的资源。
3. SetLimit、Go 与 TryGo 的准入
SetLimit(n) 限制活跃 goroutine 数,Go 在达到限制时阻塞,直到有槽位或 group 已有任务结束。它不是有界队列:循环本身停在提交位置,尚未提交的元素仍在调用者数据结构中。TryGo 在无槽位时返回 false,适合明确拒绝或改为同步处理。
在活跃 goroutine 存在时修改 limit 不符合 API 约定;启动前设置一次。限制为 0 会阻止任何任务启动,外部配置必须验证为正。若 Go 的调用者在持锁状态阻塞,而已有任务需要同一锁退出,会死锁。
group.SetLimit(maxConcurrent)
for _, task := range tasks {
task := task
if !group.TryGo(func() error { return task(ctx) }) {
return ErrOverloaded
}
}
return group.Wait()
这段提前返回时仍必须 Wait 已接纳任务;直接 return 会让它们后台继续。生产 API 应在拒绝时取消、等待,再返回 overload,或规定已接纳任务独立完成。
4. singleflight 的状态机:共享调用而非缓存
singleflight.Group.Do(key, fn) 对同一 Group、同一 key 的重叠调用只执行一个 fn,其余等待并收到相同 value/error。第三个返回值 shared 表示结果至少与其他调用共享,不表示来自缓存。调用完成后条目移除,下一次请求仍会重新执行。
type Loader struct {
group singleflight.Group
repo *Repository
}
func (l *Loader) Load(ctx context.Context, tenant, id string) (Article, error) {
key := tenant + ":" + id
value, err, _ := l.group.Do(key, func() (any, error) {
return l.repo.Find(ctx, tenant, id)
})
if err != nil {
return Article{}, fmt.Errorf("load article %q: %w", id, err)
}
article, ok := value.(Article)
if !ok {
return Article{}, errors.New("singleflight returned unexpected article type")
}
return article, nil
}
这种写法把首个调用者的 ctx 交给共享工作:首位请求取消会让所有等待者收到取消,即使他们还有预算。是否正确取决于产品语义,不能默认采用。
5. key 必须包含全部结果维度
singleflight 的 key 是隔离边界。租户、区域、语言、权限版本、字段掩码和数据版本只要影响结果,就必须进入 key。遗漏租户可能把一个租户的数据返回给另一个,是严重安全问题。key 也不能直接拼接含分隔符的任意字段,避免 a:b + c 与 a + b:c 碰撞。
可使用长度前缀、结构序列化或不可歧义编码,并限制 key 长度。不要把认证 token 原文放 key 和指标;可使用稳定、无秘密的主体/权限版本。高基数随机 key 不会受 singleflight 帮助,还可能在请求窗口内占用大量条目。
错误也会共享。短暂数据库失败到达时,所有等待者同时失败;singleflight 不做错误缓存。是否随后重试要有总预算和退避,否则一波等待者失败后会立即形成下一波回源。
6. DoChan 让等待者独立取消
Do 会同步等待共享调用,等待者无法仅靠自己的 ctx 提前离开。DoChan 返回结果 channel,等待者可 select 自己的取消。取消等待不会自动取消底层共享函数,因为其他等待者可能仍需要它。
resultCh := group.DoChan(key, func() (any, error) {
loadCtx, cancel := context.WithTimeout(context.Background(), 800*time.Millisecond)
defer cancel()
return repository.Find(loadCtx, tenant, id)
})
select {
case result := <-resultCh:
if result.Err != nil {
return Article{}, fmt.Errorf("load article: %w", result.Err)
}
article, ok := result.Val.(Article)
if !ok {
return Article{}, errors.New("unexpected article result")
}
return article, nil
case <-ctx.Done():
return Article{}, context.Cause(ctx)
}
这里共享工作使用独立、短 deadline;所有等待者离开后它最多继续 800ms。若底层工作昂贵,需要自建引用计数取消状态,但复杂度高且竞态多,通常固定短预算更可靠。
7. Forget、失效和迟到结果
Forget(key) 让后续调用不再等待当前调用,但不会停止已经运行的 fn。新旧调用可能并发,旧结果仍会返回给旧等待者。它适合明确认为旧飞行已过时的场景,不是取消 API。
缓存失效时,单纯 Forget 仍可能让慢旧加载随后写回缓存,覆盖新数据。写回必须携带 generation/version 并比较:加载开始于版本 7,期间失效推进到 8,迟到的 7 应丢弃。仅按完成时间“最后写赢”会把慢旧值当新值。
singleflight Group 通常随服务或缓存组件存在,不需 Close,也不启动常驻 goroutine;每个调用的函数资源由函数自身释放。不要每次请求创建新 Group,那样无法跨请求合并。
8. singleflight 与缓存、负缓存和回源闸门
典型读取顺序是缓存 Get,miss 后 singleflight,再在函数内部二次 Get,最后回源并 Set。二次检查防止调用排队期间其他路径已填充。缓存写失败通常只影响性能,不应把成功数据库读取变业务失败,除非缓存本身是权威存储。
L1 get -> miss -> singleflight(key)
-> L1 get again
-> acquire global semaphore
-> database/RPC
-> version-checked cache set
singleflight 只限制相同 key;一万个不同 key 仍会同时打数据库,所以回源还要全局 semaphore。不存在结果可短期负缓存,但要区分权威 not-found 与临时失败,并防止用户构造随机 key 占满缓存。
多实例各有 Group,不能跨进程合并。分布式锁会引入租约、所有者死亡和 fencing token,不应仅为“减少几次查询”贸然加入;共享缓存和下游限流常更简单。
9. semaphore 的令牌与资源所有权
semaphore.Weighted 维护总权重,Acquire(ctx, n) 等待 n 个单位,TryAcquire(n) 立即判断,Release(n) 归还。调用者一旦成功 Acquire 就拥有令牌,必须在所有返回路径 Release;失败 Acquire 不得 Release。
func transcode(ctx context.Context, sem *semaphore.Weighted, job Job) error {
weight := job.MemoryMiB
if weight <= 0 || weight > 512 {
return errors.New("job memory weight is outside allowed range")
}
if err := sem.Acquire(ctx, weight); err != nil {
return fmt.Errorf("acquire transcode capacity: %w", err)
}
defer sem.Release(weight)
if err := runTranscode(ctx, job); err != nil {
return fmt.Errorf("transcode job %q: %w", job.ID, err)
}
return nil
}
权重是估算资源量,不是任务优先级。单项权重大于总量会永远等到 context 取消,应在提交前拒绝。Release 超过已获得量属于程序错误,可能 panic,不能用 recover 当正常控制流。
10. 并发限制、速率限制和公平性
semaphore 限制同时在途,不限制每秒启动数。任务从 1 秒变 10ms 时,同样并发可产生百倍 QPS;受配额 API 还需令牌桶速率限制。数据库自身连接池已经限制连接数,应用 semaphore 应与它协调,避免两层长队列。
加权请求存在队头阻塞:一个大权重等待者可能让后续小任务无法充分利用临时空槽,具体公平行为应以版本实现和压测为准,不把它当严格调度契约。多租户公平需要每租户配额或调度器,不是一个全局 semaphore 能解决。
Acquire 等待算进请求 deadline。指标应分开队列等待和执行时间;大量 deadline 在 Acquire 处耗尽说明容量/准入错误,而不是下游执行慢。
11. 用 channel 还是 Weighted semaphore
所有任务等权时,容量为 N 的令牌 channel 简单直观,发送取得、接收归还;但每次发送都必须可取消,归还所有权要严格。Weighted 支持不同成本、context Acquire 和 TryAcquire,API 更直接。
channel 还能携带真实资源对象,形成对象池;semaphore 只代表数量,不能保证某个连接或 buffer 的生命周期。资源池要验证归还对象状态,坏连接不能放回。sync.Pool 是 GC 辅助缓存,条目可随时消失,不用于硬并发限制。
不要用 semaphore 包住一个无界 goroutine 创建循环:虽然执行被限制,等待令牌的 goroutine 数仍无界。应在启动 goroutine 前 Acquire,或让有限 worker 消费有界队列。
12. ants 的池模型与适用场景
ants v2 维护有限 worker goroutine 并复用它们执行提交函数。它适合海量、短小、彼此独立的任务,且 benchmark/profile 已证明 goroutine 创建、栈或调度成为实际瓶颈。普通 Web/RPC I/O 通常直接 goroutine 加并发上限更清晰,Go 运行时本就高效调度 goroutine。
pool, err := ants.NewPool(
64,
ants.WithExpiryDuration(30*time.Second),
ants.WithNonblocking(true),
)
if err != nil {
return fmt.Errorf("create worker pool: %w", err)
}
defer pool.Release()
if err := pool.Submit(func() {
process(job)
}); err != nil {
return fmt.Errorf("submit job: %w", err)
}
非阻塞模式在满载时返回错误,调用方必须拒绝、降级或进入其他有界持久队列。阻塞模式会把背压传给 Submit 调用者,也必须让外围请求有 deadline;具体版本的选项和错误值应以 v2.11.3 文档及契约测试为准。
13. ants 任务错误、panic 和结果回传
Submit(func()) 的函数没有 error 返回值。任务结果和错误必须通过调用者拥有的 result channel、每任务 future 或受锁结果槽传回;Submit 成功只表示已接纳,不表示任务成功。不要在任务里只日志错误然后让业务以为成功。
type taskResult struct {
value string
err error
}
resultCh := make(chan taskResult, 1)
if err := pool.Submit(func() {
value, err := process(ctx, job)
resultCh <- taskResult{value: value, err: err}
}); err != nil {
return "", fmt.Errorf("submit job: %w", err)
}
select {
case result := <-resultCh:
return result.value, result.err
case <-ctx.Done():
return "", context.Cause(ctx)
}
容量 1 使调用方取消后任务发送不会永久阻塞。任务仍须观察 ctx,否则取消只让等待者离开。panic 策略要显式配置/封装并监控;recover 后进程状态是否可信取决于任务,不能普遍转成成功。
14. Pool 生命周期、Tune 与关停
创建 pool 的组件拥有 Release 责任。先停止新提交,等待已接纳任务按业务策略完成,再 Release;仅 Release 不等于持久任务已保存。共享 pool 应在 HTTP handler 全部退出后关闭,防止关闭资源仍被使用。
动态 Tune 改容量只改变并发约束,不会创造下游容量。依据瞬时队列扩缩会振荡;应看排队分位数、下游限额、内存和冷却周期。多处调用 Tune 会丢失单一所有者。
任务闭包捕获大 request/body 会延长其生命周期直到执行完成。池虽然限制 worker,却不一定限制提交方持有的待执行对象;必须确认阻塞/队列模型和内存上限。需要重启后不丢的任务应进入持久消息系统,ants 不是作业队列。
15. 四种工具怎样安全组合
一个缓存批量加载可先用 errgroup 管理请求内多个 key,SetLimit 限制该请求并发;每个 key 用 singleflight 合并进程内重复;共享 semaphore 限制全服务数据库回源。通常不再需要 ants,机制越多越难分析等待链。
组合顺序决定资源持有。先取得 semaphore 再进入 singleflight,会让同 key 的每个等待者都占令牌,失去合并收益;应在 singleflight 的实际执行函数内 Acquire。反过来,持有稀缺连接时不要再长时间等待 pool 槽。
request errgroup limit
-> singleflight key
-> global weighted semaphore
-> cancellable repository call
<- release token
<- shared result
-> Wait all admitted request tasks
每层只承担一种容量,并共享父 deadline。画出资源获取顺序,确保所有路径反向释放,避免循环等待。
16. 超时、取消和关停协议
请求 context 只属于该请求,不要存入长期 struct。共享 singleflight 工作可用独立短预算;errgroup 任务用派生 ctx;semaphore Acquire 用当前操作 ctx;ants 任务闭包显式捕获 ctx,但池本身通常由服务生命周期管理。
关停顺序:停止入口和提交;取消允许中止的任务;等待 request groups;排空或持久化已接纳任务;Release pool;最后关闭数据库/HTTP transport。到达关停 deadline 时,进程退出会终止 goroutine 且 defer 不保证执行,所以必须持久化的工作不能只在内存。
所有阻塞点都要能回答取消来源:Go 提交是否会阻塞、Do 是否可中断、Acquire 是否带 deadline、Submit 满时怎样处理、结果发送是否可能卡住。只在任务函数开头检查一次 ctx 不够。
17. 测试并发语义而不是依赖 Sleep
用 channel barrier 控制任务开始和释放。singleflight 测试让第一个函数关闭 started,再启动多个等待者,确认执行计数为 1;取消测试用 DoChan。semaphore 测试先占满,再确认 TryAcquire false 和 Acquire 在 cancel 后返回。errgroup 测试确认错误触发其他任务观察 ctx 且 Wait 后全部退出。
func TestSingleflightSharesCall(t *testing.T) {
var group singleflight.Group
var calls atomic.Int64
started := make(chan struct{})
release := make(chan struct{})
load := func() (any, error) {
if calls.Add(1) == 1 {
close(started)
}
<-release
return "value", nil
}
first := group.DoChan("key", load)
<-started
second := group.DoChan("key", load)
close(release)
if result := <-first; result.Err != nil {
t.Fatal(result.Err)
}
if result := <-second; result.Err != nil {
t.Fatal(result.Err)
}
if got := calls.Load(); got != 1 {
t.Errorf("calls = %d, want 1", got)
}
}
测试自身加宽松 watchdog 防永久挂起,失败时输出 goroutine dump。重复和 race 扩大交错覆盖,但不通过增加 Sleep 掩盖协议问题。
18. 性能、诊断与生产指标
指标按机制区分:errgroup 活跃/首错/取消时长;singleflight 实际加载、共享等待者、key 类别与加载耗时;semaphore 获取等待、使用权重、拒绝/取消;ants running/free/cap、Submit 阻塞或拒绝、任务耗时和关停时长。不要把原始高基数 key 做标签。
goroutine profile 大量停在 semaphore Acquire 表示容量等待,大量停在 result send 表示消费者离开或缓冲错误,大量 singleflight Do 等待可能是热点回源变慢。block/mutex profile 和 trace 能还原阻塞,不要只看 CPU。
go test -count=1 ./...
go test -race ./...
go test -run TestSingleflight -count=100 ./...
go test -bench=. -benchmem ./...
go vet ./...
go tool pprof http://127.0.0.1:6060/debug/pprof/goroutine
benchmark 比较直接 goroutine、errgroup limit 与 ants 时,保持任务成本和并发度相同,记录吞吐、P95/P99、分配、峰值 goroutine 与内存。只有优势在真实负载中稳定存在才引入池。
19. 选型与上线清单
- 一批任务首错取消并等待全部退出用 errgroup;需要全部错误时显式收集,不误用首错语义。
- 同 key 同时调用合并用 singleflight;key 包含租户和全部结果维度,并结合缓存版本防迟到写。
- 资源有不同成本用 Weighted semaphore;等权任务可用简单令牌,限并发不等于限速。
- ants 只在 profile 证明大量短任务的 goroutine 成本显著时采用;池满策略、错误回传与 Release 明确。
- Acquire 在启动 goroutine 前或实际 singleflight 函数内完成,避免无界等待 goroutine和重复占令牌。
- 所有 I/O 接受 context,等待可取消;取消后仍 Wait,才能证明资源不再被使用。
- 队列、任务闭包和缓存都计入内存预算;必须持久的任务不放纯内存池。
- 固定
x/sync v0.17.0、ants/v2 v2.11.3,升级前核对 API、竞态、吞吐和关停契约。
系列导航与关联阅读
- 系列入口:Go 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go 通用工具封装:Options、分页、错误码、重试与包边界
- 下一篇:Go 依赖注入:手工装配、Uber Fx 与 Wire 的选择
- 延伸:Go 并发模式:Worker Pool、Pipeline、Fan-out 与背压
- 延伸:Go context 完整指南:取消、超时、Deadline 与 Value
- 延伸:Go Redis 生产模式:缓存一致性、穿透击穿、锁与 Streams
- 延伸:Go goroutine 生命周期:泄漏、打断、错误传播与优雅关闭
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论