conc:Go 结构化并发工具集源码级解析,及它在 inngest 中的落地实践
【免费下载链接】inngestThe leading workflow orchestration platform. Run stateful step functions and AI workflows on serverless, servers, or the edge.项目地址: https://gitcode.com/GitHub_Trending/in/inngest
conc(github.com/sourcegraph/conc,本仓库以 vendor 方式锁定在 v0.3.0)是 Sourcegraph 出品的 Go 结构化并发工具集,它以conc.WaitGroup为基石,把"goroutine 必须有主、panic 必须优雅处理、并发代码必须可读"三条原则打包进一组小而精的 API。inngest 作为一个长期运行、重度并发的编排引擎,在 pkg/service/service.go、pkg/execution/queue/backlog_normalization.go、pkg/execution/realtime/broadcaster.go 中实际使用了 conc 的WaitGroup与pool.Pool。读完本文,你将掌握 conc 的完整 API 图谱、底层实现原理(panic 捕获与传播、懒启动任务池、错误聚合),并看到它如何帮助 inngest 优雅地管理后台 goroutine 的生命周期。
什么是 conc:把结构化并发变成一件顺手的事
conc的定位是"Go 结构化并发的工具腰带"(toolbelt),它让常见的并发任务更简单、更安全。在 inngest 仓库中,它被作为直接依赖引入(见 go.mod),并随项目一起 vendor 在 vendor/github.com/sourcegraph/conc 目录下,因此其实现代码就是项目代码的一部分,可以直接阅读。
安装方式很简单:
go get github.com/sourcegraph/concconc 的使用哲学可以用三个目标概括:
- 让 goroutine 更难泄漏——所有并发必须有作用域(scoped),每个 goroutine 都必须有明确的所有者;
- 优雅处理 panic——子 goroutine 的 panic 会被捕获、附带堆栈信息并传播给等待者,而不是直接打崩整个进程;
- 让并发代码更容易阅读——用高层抽象抹平繁琐的样板代码。
一图速览:conc 的完整 API 地图
原版 README 用一张速查表总结了 conc 各子包的能力,下面完整保留并补充了使用场景:
| 你需要的能力 | 使用哪个 API | 典型场景 |
|---|---|---|
更安全的sync.WaitGroup | conc.WaitGroup | 起一组 goroutine 并等待,自动处理 panic |
| 并发受限的任务执行器 | pool.Pool | 最多 N 个 goroutine 并发执行任务 |
| 并发执行并收集任务结果 | pool.ResultPool | 任务有返回值(Go 泛型) |
| 任务可能失败 | pool.ErrorPool | 任务返回 error,Wait()聚合错误 |
| 失败时任务应被取消 | pool.ContextPool | 任务共享一个 context,出错即整体取消 |
| 有序流的并行处理 | stream.Stream(stream 子包) | 并行处理数据流但回调保持提交顺序 |
| 并发 map 一个 slice | iter.Map(iter 子包) | input映射为output |
| 并发遍历一个 slice | iter.ForEach(iter 子包) | 对每个元素并发执行副作用操作 |
| 在自己的 goroutine 里捕获 panic | panics.Catcher | 手动管理 panic 的捕获与恢复 |
所有任务池都由pool.New()(无结果)或pool.NewWithResults[T]()(有结果)创建,然后通过链式方法配置:
p.WithMaxGoroutines(n):限制池内最大 goroutine 数,默认不限制,n < 1时直接 panic;p.WithErrors():把池升级为可运行返回 error 任务的 ErrorPool;p.WithContext(ctx):让池内任务共享 context,可通过WithCancelOnError()让首个错误触发整体取消;p.WithFirstError():错误池只保留第一个错误,而不是聚合错误;p.WithCollectErrored():结果池即使任务出错也收集其结果(默认出错任务的结果被丢弃)。
一个关键约束值得注意:所有With*配置方法一旦在第一次调用Go()之后被调用,就会 panic(源码中通过panicIfInitialized()强制保证),因为任务池在启动后不允许再被重新配置。这在 pool.go 中有明确实现。
目标一:让 goroutine 更难泄漏——作用域并发的 WaitGroup
用go关键字随手起 goroutine 最大的痛点之一是清理:非常容易"发出去"却忘了"等回来",导致 goroutine 泄漏。conc 对此采取了一个强立场:所有并发都应该是作用域化的——goroutine 必须有一个所有者,所有者必须保证自己拥有的 goroutine 正确退出。
在 conc 中,goroutine 的所有者永远是conc.WaitGroup。goroutine 通过(*WaitGroup).Go()产生,并且在 WaitGroup 离开作用域之前必须调用(*WaitGroup).Wait()。
WaitGroup的实现非常朴素,但它把"等待"和"panic 收集"绑定在了一起(见 waitgroup.go):
type WaitGroup struct { wg sync.WaitGroup pc panics.Catcher } func (h *WaitGroup) Go(f func()) { h.wg.Add(1) go func() { defer h.wg.Done() h.pc.Try(f) }() } func (h *WaitGroup) Wait() { h.wg.Wait() // Propagate a panic if we caught one from a child goroutine. h.pc.Repanic() } func (h *WaitGroup) WaitAndRecover() *panics.Recovered { h.wg.Wait() return h.pc.Recovered() }可以看到:WaitGroup的内部只是sync.WaitGroup加上一个panics.Catcher,零值可用、用法与标准库一致;区别在于Go()用pc.Try(f)包住了每个任务,Wait()在等待结束后会Repanic()把子 goroutine 的 panic 重新抛给调用者。
如果你的 goroutine 需要比调用者活得更久,可以把WaitGroup作为参数传进产生 goroutine 的函数:
func main() { var wg conc.WaitGroup defer wg.Wait() startTheThing(&wg) } func startTheThing(wg *conc.WaitGroup) { wg.Go(func() { ... }) }关于"go 语句是有害的"以及作用域并发为什么更优雅的讨论,可参考结构化并发的经典论述(vorpus.org 的Notes on structured concurrency),conc 正是把这一思想做成了开箱即用的库。
目标二:优雅处理 panic——捕获、装饰、再传播
长期运行的应用中,一个没有 panic handler 的 goroutine 一旦 panic 会拖垮整个进程。但如果自己加 handler,捕获之后怎么办?常见选择有四种:忽略、打日志、转成 error 返回给 spawner、把 panic 传播给 spawner。
conc 的判断是:忽略是坏主意(panic 通常意味着真的有 bug);只打日志也不好(spawner 得不到任何信号,程序可能带病继续跑)。合理的做法是 (3)(4),但它们都要求 goroutine 有一个能真正接收"出事了"消息的所有者——普通go语句做不到,而 conc 的所有 goroutine 都有所有者。
因此在 conc 里:任何一次Wait()调用,只要子 goroutine panic 过,就会以该 panic 值重新 panic,并且 panic 值会被装饰上子 goroutine 的堆栈信息(debug.Stack()),这样你不会丢失事故现场。
panics.Catcher 的源码实现
Catcher的核心是一个atomic.Pointer[Recovered],因此它对多个 goroutine 并发调用Try是安全的,且只保留第一个捕获到的 panic(见 panics.go):
type Catcher struct { recovered atomic.Pointer[Recovered] } func (p *Catcher) Try(f func()) { defer p.tryRecover() f() } func (p *Catcher) tryRecover() { if val := recover(); val != nil { rp := NewRecovered(1, val) p.recovered.CompareAndSwap(nil, &rp) } } func (p *Catcher) Repanic() { if val := p.Recovered(); val != nil { panic(val) } }Recovered结构体携带三类信息(panics.go):
Value any:panic 的原始值;Callers []uintptr:runtime.Callers记录的调用栈 PC,可用runtime.CallersFrames还原更详细的栈帧;Stack []byte:捕获时debug.Stack()得到的格式化堆栈,开箱即用。
它还提供两个实用转换:String()输出人可读的panic: ...\nstacktrace:\n...格式;AsError()把 panic 转成 error(ErrRecovered),并且实现Unwrap()——如果 panic 值本身是 error,可以被errors.Is/errors.As解开。另外,panics.Try(f)是一个独立工具函数:执行f并返回捕获到的*Recovered,你可以选择panic()重新传播,或者AsError()当普通错误处理(见 try.go)。
对比:标准库手写 vs conc
下面这张来自 README 的对照表最能说明问题。左边是标准库手写"捕获 panic → 记录堆栈 → 通过 channel 传回 → 再 panic"的全过程,右边是 conc 的全部代码:
| stdlib | conc |
|---|---|
| ```go |
type caughtPanicError struct { val any stack []byte }
func (e *caughtPanicError) Error() string { return fmt.Sprintf( "panic: %q\n%s", e.val, string(e.stack) ) }
func main() { done := make(chan error) go func() { defer func() { if v := recover(); v != nil { done <- &caughtPanicError{ val: v, stack: debug.Stack() } } else { done <- nil } }() doSomethingThatMightPanic() }() err := <-done if err != nil { panic(err) } }|go func main() { var wg conc.WaitGroup wg.Go(doSomethingThatMightPanic) // panics with a nice stacktrace wg.Wait() }
每次用 `go` 手动做完这一整套都相当繁琐,而且样板代码会淹没业务逻辑的可读性——这正是 conc 替你完成的事情。 ## 目标三:让并发代码更容易阅读——五个高频场景对照 正确写并发很难,写得既不绕又能让人一眼看懂更难。conc 用高层抽象抹平样板代码,下面五个场景(均来自 README,为简洁起见省略了 panic 传播)展示了标准库写法与 conc 写法的差距。 ### 场景 1:起一组 goroutine 并等待 | stdlib | conc | | --- | --- | | ```go func main() { var wg sync.WaitGroup for i := 0; i < 10; i++ { wg.Add(1) go func() { defer wg.Done() // crashes on panic! doSomething() }() } wg.Wait() } ``` | ```go func main() { var wg conc.WaitGroup for i := 0; i < 10; i++ { wg.Go(doSomething) } wg.Wait() } ``` | 标准库版本不但要手动 `Add`/`Done`,而且任何子 goroutine panic 都会直接崩溃整个进程;conc 版本自动处理了这一切。 ### 场景 2:在固定大小的 goroutine 池中处理一个流 | stdlib | conc | | --- | --- | | ```go func process(stream chan int) { var wg sync.WaitGroup for i := 0; i < 10; i++ { wg.Add(1) go func() { defer wg.Done() for elem := range stream { handle(elem) } }() } wg.Wait() } ``` | ```go func process(stream chan int) { p := pool.New().WithMaxGoroutines(10) for elem := range stream { elem := elem p.Go(func() { handle(elem) }) } p.Wait() } ``` | ### 场景 3:在固定大小 goroutine 池中处理一个 slice | stdlib | conc | | --- | --- | | ```go func process(values []int) { feeder := make(chan int, 8) var wg sync.WaitGroup for i := 0; i < 10; i++ { wg.Add(1) go func() { defer wg.Done() for elem := range feeder { handle(elem) } }() } for _, value := range values { feeder <- value } close(feeder) wg.Wait() } ``` | ```go func process(values []int) { iter.ForEach(values, handle) } ``` | 一个手动构建 feeder channel + worker 池的经典模式,被压缩成一行 `iter.ForEach`。 ### 场景 4:并发 map 一个 slice | stdlib | conc | | --- | --- | | ```go func concMap( input []int, f func(int) int, ) []int { res := make([]int, len(input)) var idx atomic.Int64 var wg sync.WaitGroup for i := 0; i < 10; i++ { wg.Add(1) go func() { defer wg.Done() for { i := int(idx.Add(1) - 1) if i >= len(input) { return } res[i] = f(input[i]) } }() } wg.Wait() return res } ``` | ```go func concMap( input []int, f func(*int) int, ) []int { return iter.Map(input, f) } ``` | 标准库方案需要手写原子计数器来分配任务下标、并发写结果切片;`iter.Map` 一行搞定(注意 conc 的 `Map` 回调签名是 `func(*T) T`,直接修改原 slice 元素)。 ### 场景 5:保持顺序的并行流处理 | stdlib | conc | | --- | --- | | ```go func mapStream( in chan int, out chan int, f func(int) int, ) { tasks := make(chan func()) taskResults := make(chan chan int) // Worker goroutines var workerWg sync.WaitGroup for i := 0; i < 10; i++ { workerWg.Add(1) go func() { defer workerWg.Done() for task := range tasks { task() } }() } // Ordered reader goroutines var readerWg sync.WaitGroup readerWg.Add(1) go func() { defer readerWg.Done() for result := range taskResults { item := <-result out <- item } }() // Feed the workers with tasks for elem := range in { resultCh := make(chan int, 1) taskResults <- resultCh tasks <- func() { resultCh <- f(elem) } } // We've exhausted input. // Wait for everything to finish close(tasks) workerWg.Wait() close(taskResults) readerWg.Wait() } ``` | ```go func mapStream( in chan int, out chan int, f func(int) int, ) { s := stream.New().WithMaxGoroutines(10) for elem := range in { elem := elem s.Go(func() stream.Callback { res := f(elem) return func() { out <- res } }) } s.Wait() } ``` | `stream.Stream` 的设计是:任务并发执行,但每个任务返回一个 `stream.Callback`,回调按**提交顺序**串行执行——这样既拿到并行计算的速度,又保住了输出的顺序。标准库实现则需要双 channel(tasks + taskResults)+ 两组 WaitGroup 才能达到同样效果。 ## 任务池家族:Pool、ResultPool、ErrorPool、ContextPool 及其组合 conc 的任务池是本文最有工程价值的部分,也是 inngest 实际用到的组件。六个池子由基础 `Pool` 通过 `With*` 方法线性组合出来:Pool ──WithErrors()──▶ ErrorPool ──WithContext(ctx)──▶ ContextPool │ │ │ └──WithFirstError()── 只保留第一个错误 │ └──WithContext(ctx)──▶ ContextPool ──WithCancelOnError()── 出错即取消
NewWithResultsT ──▶ ResultPool[T] │ │ │──WithErrors()──▶ ResultErrorPool[T] ──WithCollectErrored()── 出错也收集结果 │ │ └──WithContext(ctx)──▶ ResultContextPool[T]
### 基础 Pool:懒启动的 goroutine 调度器 `Pool` 的结构体([pool.go](https://link.gitcode.com/i/f9ad865dd50f5bc534a6c50fd7b5a3f5#L30-L35))只有四样东西:一个内部 `conc.WaitGroup`(负责等待与 panic 传播)、一个 `limiter`(即 `chan struct{}`,容量即最大 goroutine 数)、一个无缓冲 `tasks chan func()`、以及一个 `initOnce`。 几个关键实现细节: - **零值可用、创建廉价**:`New()` 只是返回空结构体,`tasks` channel 在第一次 `Go()`/`Wait()` 时由 `initOnce.Do` 惰性初始化([pool.go](https://link.gitcode.com/i/f9ad865dd50f5bc534a6c50fd7b5a3f5#L100-L104)); - **goroutine 数永远不会超过任务数**:`Go()` 会先尝试把任务塞给空闲 worker,只有塞不进去时才启动新 worker([pool.go](https://link.gitcode.com/i/f9ad865dd50f5bc534a6c50fd7b5a3f5#L39-L70)); - **效率有边界**:注释明确说明 Pool 不是零成本——启动/收尾约 1µs、每个任务约 300ns 开销,不适合超短任务; - **`WithMaxGoroutines(n)` 中 `n < 1` 直接 panic**([pool.go](https://link.gitcode.com/i/f9ad865dd50f5bc534a6c50fd7b5a3f5#L87-L96))。 ### ErrorPool:错误收集与聚合 `ErrorPool.Go(f func() error)` 把任务包进基础池,用互斥锁保护错误收集([error_pool.go](https://link.gitcode.com/i/54302d08f4129a621a6e523b4916fa82#L85-L97))。默认 `Wait()` 返回的是用 `multierror.Join` 聚合起来的**组合错误**(见 [internal/multierror](https://link.gitcode.com/i/74b7764a87d69c91ac5197f33477b53d),按 Go 1.20+ 的 `errors.Join` 语义实现);调用 `WithFirstError()` 后只保留第一个错误。 ### ContextPool:共享取消 `WithContext(ctx)` 会在内部对传入 ctx 再包一层 `context.WithCancel`([pool.go](https://link.gitcode.com/i/f9ad865dd50f5bc534a6c50fd7b5a3f5#L138-L146)),任务拿到的是这个可被池取消的子 ctx。配置 `WithCancelOnError()` 后,**只要任意任务返回 error 或 panic,池就会 cancel 掉共享 context**,其余任务随即感知取消([context_pool.go](https://link.gitcode.com/i/e73370eb2285920521b2eb33e1451cd5#L24-L50))。 实现上有个值得注意的细节:取消时当前错误是直接通过 `addErr` 写入的(绕过了 ErrorPool 包装),这是为了避免"取消导致其他 goroutine 先返回 `context.Canceled`,反而抢占了 `WithFirstError()` 的第一个错误位"。也因此官方建议 `WithCancelOnError()` 与 `WithFirstError()` 搭配使用——第一个错误之后的所有错误大概率都是 `context.Canceled`。 ### ResultPool 系列:泛型结果收集 `pool.NewWithResults[T]()` 创建 `ResultPool[T]`,任务签名 `func() T`,`Wait()` 返回 `[]T`([result_pool.go](https://link.gitcode.com/i/97bb4ce1ea47f893a23618dccfc5b64a#L25-L43))。结果收集用 `resultAggregator[T]`(互斥锁 + append)完成。**注意结果顺序不保证与提交顺序一致**——需要顺序请用 `stream` 或 `iter.Map`。 `ResultErrorPool[T]`(任务 `func() (T, error)`)默认丢弃出错任务的结果,`WithCollectErrored()` 改为照常收集([result_error_pool.go](https://link.gitcode.com/i/a7f5cbbdf3d826829f60f0748176ce0d));`ResultContextPool[T]`(任务 `func(context.Context) (T, error)`)同理叠加 context 取消能力([result_context_pool.go](https://link.gitcode.com/i/05b116e92daee1ec7fc23abe69635402))。 ## 在 inngest 中的真实落地:三个生产级用法 inngest 的源码给了 conc 三个教科书级的用法示范,这也印证了 README 中"goroutine 必须有所有者、panic 必须优雅处理"的设计目标。 ### 1. 全局 WaitGroup:所有后台 goroutine 的唯一所有者 [pkg/service/service.go](https://link.gitcode.com/i/1043ff4a2c47d3b7367705fd7dd99c3a#L23-L31) 定义了一个包级 `conc.WaitGroup`,并导出 `Go`/`Wait` 两个薄封装,全服务的后台 goroutine 都从这里产生: ```go var wg conc.WaitGroup func Go(f func()) { wg.Go(f) } func Wait() { wg.Wait() }而在优雅停机流程里(pkg/service/service.go),服务 Stop 之后会调用wg.WaitAndRecover()——注意这里刻意不用会 re-panic 的Wait(),而是取出*panics.Recovered后把 panic 值和堆栈记入日志,避免停机过程本身被 panic 打断:
if recovered := wg.WaitAndRecover(); recovered != nil { l.Error("global goroutine panic waiting for service to stop", "error", recovered.Value, "stack", recovered.Stack) }这正是 concWaitGroup提供WaitAndRecover这一姊妹方法的动机:等待与 panic 传播解耦,让调用者自行决定"重新 panic"还是"降级为日志"。
2. 受限并发池:backlog 归一化
pkg/execution/queue/backlog_normalization.go 用pool.New().WithMaxGoroutines(...)限制 backlog 归一化的并发度,然后逐条提交NormalizeItem任务并统一wg.Wait():
wg := pool.New().WithMaxGoroutines(int(q.backlogNormalizeConcurrency)) for _, item := range res.Items { item := item // capture range variable wg.Go(func() { _, err := q.NormalizeItem(logger.WithStdlib(ctx, l), sp, latestConstraints, backlog, *item) if err != nil && !errors.Is(err, context.Canceled) { l.ReportError(err, "could not normalize item", ...) } }) } wg.Wait()这里体现了WithMaxGoroutines的实战价值:队列归一化可能面对海量 backlog 条目,必须用显式并发上限保护数据库,同时pool.Pool自带的 panic 传播保证任何一个归一化任务的异常都不会静默吞掉。
3. 广播器的等待与恢复
pkg/execution/realtime/broadcaster.go 用wg conc.WaitGroup管理各runTopicgoroutine,并在收尾时用WaitAndRecover()检查是否有子 goroutine panic(broadcaster.go),与 service.go 的模式如出一辙:实时广播这种长期运行、不可中断的路径上,panic 一律捕获为日志而不是炸掉整个进程。
版本状态与注意点
本仓库 vendor 的是github.com/sourcegraph/conc v0.3.0(见 go.mod),属于 pre-1.0 阶段。README 中官方声明:1.0 之前 API 可能仍有小规模破坏性变更(主要是稳定 API 和调整默认值),因此升级依赖时需留意 release notes。对 inngest 这样的生产项目而言,将其 vendor 进仓库意味着 API 变更不会悄悄发生——这是把第三方并发库纳入版本控制的一个务实做法。
小结:什么时候用 conc,什么时候不用
| 场景 | 推荐方案 |
|---|---|
| 只想起一组 goroutine 并安全等待 | conc.WaitGroup(替代sync.WaitGroup) |
| 需要并发上限 + panic 安全 + 错误收集 | pool.New().WithErrors().WithMaxGoroutines(n) |
| 需要首错即取消整体 | pool.New().WithContext(ctx).WithCancelOnError()(建议搭配WithFirstError()) |
| 并发处理 slice/流且需保序 | iter.Map/iter.ForEach/stream.Stream |
| 自己的 goroutine 里想捕获 panic 再决定去向 | panics.Catcher或panics.Try |
| 任务极短(微秒级以下)、追求极致零开销 | 谨慎评估:Pool 每个任务约 300ns 开销 |
conc 的核心价值不在性能(它自己也承认不是零成本),而在于把"并发必须有所有者、panic 必须可追踪、代码必须可读"这三条纪律变成 API 的默认行为。inngest 在服务生命周期、队列处理、实时广播三处关键路径上的用法,正是这三条纪律在生产环境中的完整样本。
【免费下载链接】inngestThe leading workflow orchestration platform. Run stateful step functions and AI workflows on serverless, servers, or the edge.项目地址: https://gitcode.com/GitHub_Trending/in/inngest
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考