safego
为什么要有它
实际业务里经常要「并行做几件互不依赖的事」——比如初始化时同时连数据库、拉配置、预热缓存。最朴素的做法是开一堆 goroutine,自己用 sync.WaitGroup + 一个 error channel 收集错误,但样板代码又长又容易漏(忘了 Wait、忘了 close channel、忘了把错误传出来)。
safego 用标准扩展库 golang.org/x/sync/errgroup 把这套模式封装成两个函数:传入一组 WorkFunc,并发跑,等全部结束,返回第一个 error 或 nil。
代码
import (
"context"
"golang.org/x/sync/errgroup"
)
// WorkFunc is simple work func
type WorkFunc func() error
// WorkFuncWithCtx is simple work func
type WorkFuncWithCtx func(ctx context.Context) error
// MultiRun run all WorkFunc, it will return first err or nil.
// caution: It waits all WorkFunc done, err won't interrupt or notice other WorkFunc when it happened.
func MultiRun(fs ...WorkFunc) error {
if len(fs) == 0 {
return nil
}
eg := &errgroup.Group{}
for i := range fs {
if fs[i] == nil {
continue
}
eg.Go(fs[i])
}
return eg.Wait()
}
// MultiRunWithCtx is mostly like MultiRun, but with ctx notify.
func MultiRunWithCtx(fs ...WorkFuncWithCtx) error {
if len(fs) == 0 {
return nil
}
eg, ctx := errgroup.WithContext(context.Background())
for i := range fs {
if fs[i] == nil {
continue
}
j := i
eg.Go(func() error {
return fs[j](ctx)
})
}
return eg.Wait()
}
MultiRun:并行跑,互不打断
MultiRun 用 errgroup.Group{}(零值即可用)把每个 WorkFunc 丢进 eg.Go 并发执行:
eg.Go(fs[i])直接把函数当参数传进去——WorkFunc本身就是func() error,和eg.Go要的签名一致,所以不需要再包一层闭包,也就没有「循环变量捕获」的坑。- 遇到
nil的函数直接continue跳过,调用方不用怕传空。 eg.Wait()会阻塞到所有函数都跑完,然后返回第一个非nil的 error(其余 error 被丢弃)。
关键点(对应注释里的 caution):MultiRun 没有 context,某个函数出错时,errgroup 既不会取消任何东西,也不会通知其他函数。其他函数会一直跑到自己结束。所以如果你的 WorkFunc 里在跑长任务,它不会因为「兄弟任务失败了」而提前停下——它根本不知道。
MultiRunWithCtx:带 ctx 通知
MultiRunWithCtx 行为大体一样,但多了一层 context:
errgroup.WithContext返回一个派生的ctx,以及绑定到它身上的errgroup.Group。- 只要任意一个
WorkFunc返回非nilerror(或外部取消这个 ctx),派生 ctx 会被自动取消。其他函数手里拿到的ctx就会变成Done,它们可以主动select ctx.Done()退出来「提前止损」。 - 因为
WorkFuncWithCtx需要 ctx 入参,eg.Go要的却是无参func() error,所以这里必须包一层闭包func() error { return fs[j](ctx) }。注意j := i这个捕获——把下标拷到局部变量再闭包,避免所有闭包共享同一个循环变量i。
但要清醒:WithContext 只是「通知」,不会物理抢占正在跑的 goroutine。errgroup 的 Wait 仍然会等所有函数返回。能不能早停,完全看你的函数有没有去监听 ctx.Done()。
两个版本怎么选
| 场景 | 用哪个 |
|---|---|
| 几个独立任务并行,谁失败都无所谓,等全部跑完拿第一个错 | MultiRun |
| 任务里有长耗时的循环 / 网络调用,希望一个失败就让其他任务尽快退出 | MultiRunWithCtx,并且在函数体里监听 ctx.Done() |
| 任务之间需要共享数据或顺序依赖 | 都不合适,回到 channel / 普通 goroutine 自己编排 |
两者的共同前提:传入的函数本身应当是并发安全的,它们共享的变量要么各自独立,要么用锁 / channel 保护。safego 只负责「并发调度 + 错误聚合」,不替你管数据竞争。
用法示例
// 真实项目中 safego 一般放在 helper 包,调用时带 helper. 前缀(与本章前面的定义是同一套 API)。
// 按位 flag 决定并行填充哪些用户信息:需要的才加进切片,最后一次性并发跑,任一失败返回第一个 error。
func covertUserInfo(flag int, userIds []int64, titleId int64, respData *Resp) error {
fs := make([]helper.WorkFunc, 0)
if flag&CovertUserInfoFlagUser != 0 {
fs = append(fs, func() error { return fillUserBasic(userIds, titleId, respData) })
}
if flag&CovertUserInfoFlagWallet != 0 {
fs = append(fs, func() error { return fillUserWallet(userIds, respData) })
}
if flag&CovertUserInfoFlagLevel != 0 {
fs = append(fs, func() error { return fillUserLevel(userIds, respData) })
}
if flag&CovertUserInfoFlagConnect != 0 {
fs = append(fs, func() error { return fillUserConnect(userIds, respData) })
}
if flag&CovertUserInfoFlagSource != 0 {
fs = append(fs, func() error { return fillUserSource(respData) })
}
if flag&CovertUserInfoFlagOik != 0 {
fs = append(fs, func() error { return fillUserOik(userIds, respData) })
}
return helper.MultiRun(fs...)
}
// 带 ctx 的版本:某个任务失败时,longTask 能通过 ctx 提前退出
func initAll() error {
return helper.MultiRunWithCtx(
func(ctx context.Context) error { return longTask(ctx) },
func(ctx context.Context) error { return fetchRemote(ctx) },
)
}
func longTask(ctx context.Context) error {
for {
select {
case <-ctx.Done(): // 兄弟任务失败 → ctx 取消 → 这里立刻退出
return ctx.Err()
default:
// 干活……
}
}
}
安装依赖:
go get golang.org/x/sync/errgroup。更多errgroup细节(如SetLimit限制并发数)见官方文档。