Go 并发任务池(fan-out / fan-in)
难度:⭐⭐ 中等
考点
- goroutine + channel 并发模型
- fan-out(分发)与 fan-in(合并)
- channel 关闭语义(生产者 close、
range退出) - context 取消与超时
题目描述
实现 RunPool:把 jobs 分发给 workers 个 worker 并发处理(fan-out),汇总结果返回(fan-in)。
要求:
workers个 goroutine 同时消费任务,处理函数fn(ctx, job)返回(int, error)- 所有任务处理完成后返回结果切片(顺序不做要求)
ctx取消时尽快返回ctx.Err(),不要阻塞fn返回 error 时,整个池返回该 error
函数签名
go
type JobFunc func(ctx context.Context, job int) (int, error)
func RunPool(ctx context.Context, jobs []int, workers int, fn JobFunc) ([]int, error)示例
go
results, err := RunPool(ctx, []int{1, 2, 3}, 3, func(ctx context.Context, j int) (int, error) {
return j * j, nil
})
// results 包含 {1, 4, 9}(顺序不限),err == nil提示
- 用两个 channel:
in(任务输入)、out(结果输出) - 关闭语义:
in由生产者close;out必须等所有 worker 结束(sync.WaitGroup)后再close,否则发送端 panic - 取消:生产与消费两端都用
select监听ctx.Done() - 思考:如果某个 worker 出错,如何让整个池尽快停止?(错误 channel + 广播退出)
参考答案(Go)
点击展开参考答案
go
//go:build ignore
package answer
import (
"context"
"sync"
)
type JobFunc func(ctx context.Context, job int) (int, error)
// RunPool 参考答案:fan-out + fan-in + context 取消
func RunPool(ctx context.Context, jobs []int, workers int, fn JobFunc) ([]int, error) {
if workers <= 0 {
workers = 1
}
in := make(chan int)
out := make(chan int)
errCh := make(chan error, workers)
var wg sync.WaitGroup
for w := 0; w < workers; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for j := range in {
v, err := fn(ctx, j)
if err != nil {
errCh <- err
return
}
select {
case out <- v:
case <-ctx.Done():
return
}
}
}()
}
go func() {
wg.Wait()
close(out)
}()
go func() {
defer close(in)
for _, j := range jobs {
select {
case in <- j:
case <-ctx.Done():
return
}
}
}()
var results []int
for v := range out {
results = append(results, v)
}
if ctx.Err() != nil {
return nil, ctx.Err()
}
select {
case err := <-errCh:
return nil, err
default:
}
return results, nil
}