在扇出/扇入工作池中,上下文取消后出现goroutine泄漏
我有一个扇出/扇入的工作池,本应在上下文被取消时停止所有工作线程,但在高负载、短超时的情况下,goroutine数量持续攀升,且从未恢复。
func processJobs(ctx context.Context, jobs []Job) ([]Result, error) {
results := make(chan Result) // unbuffered
var wg sync.WaitGroup
for _, job := range jobs {
wg.Add(1)
go func(j Job) {
defer wg.Done()
select {
case results <- doWork(j):
case <-ctx.Done():
return
}
}(job)
}
go func() {
wg.Wait()
close(results)
}()
var out []Result
for {
select {
case r, ok := <-results:
if !ok {
return out, nil
}
out = append(out, r)
case <-ctx.Done():
return out, ctx.Err()
}
}
}
我原以为在每个工作者的 select 中的 ctx.Done() 情况,一旦上下文被取消就会立即退出。相反,在带有短期限的负载下,goroutine的数量持续增长,即使请求完成后也不会回落。
到底是什么原因导致这里的泄漏,以及正确的修复方案是什么?
解决方案
doWork(j)不会被上下文取消,因此始终可能阻塞。doWork也需要对上下文做出响应,才能被取消。
case results <- doWork(j): 并不意味着“在上下文先被取消时才发送 doWork(j) 的结果”。
对于一个 select,发送语句的通道操作数和右值表达式在选取case之前就会被求值。所以 doWork(j) 先执行。goroutine会尝试发送结果,或者注意到 ctx.Done() 已就绪,……只有在 doWork 完成之后才会发生。
这意味着取消并不会中止 doWork。它只会中止最后一次发送。
你可以通过这个 小型Playground示例 看到这一点。
输出:
worker: entering select
slowWork: started; it does not observe ctx
main: canceling context
main: no result yet; caller returns
slowWork: finished
worker: noticed cancellation, but only after slowWork returned
main: exit
在这里,slowWork 仍然在上下文被取消后完成执行。select 只能在发送的值已经被计算后才对取消做出响应。
让工作本身具备上下文感知能力,避免在接收方已经返回时被无限阻塞:
r, err := doWork(ctx, j)
if err != nil {
return
}
select {
case results <- r:
case <-ctx.Done():
return
}
把 ctx 传给 doWork 让工作在上下文被取消时就提前停止。
然后,环绕 results <- r 的 select 可以防止在调用方已经返回时发送结果而让工作者卡住。
在 playground 的第二个示例会产生以下输出:
worker 1: doWork started
worker 2: doWork started
main: processJobs returned out=[] err=context deadline exceeded
worker 1: stopped: context deadline exceeded
worker 2: stopped: context deadline exceeded
main: exit
在这里,doWork 本身会感知上下文并在被取消时返回。若没有它,在它正在运行时外部的 select 就无法对其进行取消。