Go 高并发实战:Worker Pool、限流、背压与优雅退出
Go 可以轻松启动 Goroutine,但“能并发”并不等于“能稳定承受高并发”。生产系统还要控制任务数量、内存占用、下游压力、超时传播和停机过程。
本文从一个原则出发:系统中的每个队列和并发点都必须有上限,每个后台任务都必须能退出。
先定义容量而不是先启动 Goroutine
以下写法会为每个任务启动一个 Goroutine:
for _, job := range jobs {
go handle(job)
}
任务数量较小时没有问题,但面对突发流量可能产生以下后果:
- Goroutine 数量持续增长,占用栈和调度资源。
- 所有任务同时请求数据库或下游服务,迅速耗尽连接池。
- 任务在内存中排队,延迟和内存同时上升。
- 服务退出时无法确认哪些任务已经完成。
设计并发系统前,应明确四个容量参数:
- 同时执行多少任务。
- 队列最多容纳多少任务。
- 每个任务允许执行多久。
- 下游每秒最多接受多少请求。
Worker Pool:限制同时执行的任务数
使用固定数量的 Worker 消费共享任务通道:
package worker
import (
"context"
"fmt"
"golang.org/x/sync/errgroup"
)
type Job struct {
ID int64
Payload []byte
}
type Handler func(context.Context, Job) error
func Run(ctx context.Context, workers int, jobs <-chan Job, handle Handler) error {
if workers < 1 {
return fmt.Errorf("workers must be positive")
}
group, ctx := errgroup.WithContext(ctx)
for i := 0; i < workers; i++ {
group.Go(func() error {
for {
select {
case <-ctx.Done():
return ctx.Err()
case job, ok := <-jobs:
if !ok {
return nil
}
if err := handle(ctx, job); err != nil {
return fmt.Errorf("handle job %d: %w", job.ID, err)
}
}
}
})
}
return group.Wait()
}
errgroup.WithContext 会在任一 Worker 返回错误时取消其他 Worker。任务生产方负责关闭 jobs 通道,消费方不能关闭自己没有创建的通道。
调用示例:
jobs := make(chan worker.Job, 100)
go func() {
defer close(jobs)
for _, job := range input {
select {
case jobs <- job:
case <-ctx.Done():
return
}
}
}()
if err := worker.Run(ctx, 8, jobs, process); err != nil {
return err
}
Worker 数量不是越大越好。CPU 密集任务通常接近 CPU 核数;I/O 密集任务可以更高,但仍受数据库连接池、远程服务限额和内存约束。
信号量:限制局部并发
不需要完整任务队列时,可以用带缓冲通道作为信号量:
func ForEach(ctx context.Context, items []Item, limit int) error {
group, ctx := errgroup.WithContext(ctx)
semaphore := make(chan struct{}, limit)
for _, item := range items {
item := item
select {
case semaphore <- struct{}{}:
case <-ctx.Done():
return ctx.Err()
}
group.Go(func() error {
defer func() { <-semaphore }()
return process(ctx, item)
})
}
return group.Wait()
}
获取信号量时也必须监听 Context,否则取消后主循环可能仍被阻塞。
对于复杂场景,可以使用 golang.org/x/sync/semaphore 的加权信号量,让不同任务占用不同权重。
背压:队列满时必须做出决定
有界队列满了以后,系统只有几种选择:
- 等待空位,让上游自然减速。
- 在限定时间内等待,超时后返回失败。
- 立即拒绝并让调用方稍后重试。
- 丢弃允许丢失的低优先级任务。
- 将任务持久化到外部消息队列。
带超时的入队函数:
var ErrQueueFull = errors.New("queue is full")
func Enqueue(ctx context.Context, queue chan<- Job, job Job, timeout time.Duration) error {
timer := time.NewTimer(timeout)
defer timer.Stop()
select {
case queue <- job:
return nil
case <-timer.C:
return ErrQueueFull
case <-ctx.Done():
return ctx.Err()
}
}
不要默默创建更大的内存队列。更大的队列只是把过载转换为更长延迟,最终仍可能耗尽内存。
HTTP 服务在过载时可以返回 429 Too Many Requests 或 503 Service Unavailable,并在适当时提供 Retry-After。
限流:约束单位时间内的请求数
golang.org/x/time/rate 实现令牌桶限流:
type Limiter struct {
limiter *rate.Limiter
}
func NewLimiter(requestsPerSecond float64, burst int) *Limiter {
return &Limiter{
limiter: rate.NewLimiter(rate.Limit(requestsPerSecond), burst),
}
}
func (l *Limiter) Middleware(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !l.limiter.Allow() {
w.Header().Set("Retry-After", "1")
http.Error(w, "too many requests", http.StatusTooManyRequests)
return
}
next.ServeHTTP(w, r)
})
}
单进程限流只约束当前实例。多副本服务需要基于网关、Redis 或专用限流组件实现全局策略。
限流维度应与业务目标匹配,例如用户、租户、API Key、来源 IP 或具体接口。只使用来源 IP 可能误伤共享出口网络。
合并重复请求
缓存失效时,大量请求可能同时加载同一个键。singleflight 可以让同一进程中相同键的并发调用共享一次结果:
type Loader struct {
group singleflight.Group
store Store
}
func (l *Loader) Load(ctx context.Context, key string) (Value, error) {
result, err, _ := l.group.Do(key, func() (any, error) {
return l.store.Load(ctx, key)
})
if err != nil {
return Value{}, err
}
return result.(Value), nil
}
注意:首次请求的执行时间会影响等待它的其他请求。加载函数仍需设置超时,跨实例场景还要依赖分布式缓存或协调机制。
批处理:减少下游往返次数
写数据库、发送日志或调用批量 API 时,可以按数量或时间窗口聚合任务:
func Batch(ctx context.Context, input <-chan Item, maxSize int, interval time.Duration, flush func([]Item) error) error {
ticker := time.NewTicker(interval)
defer ticker.Stop()
items := make([]Item, 0, maxSize)
flushItems := func() error {
if len(items) == 0 {
return nil
}
batch := append([]Item(nil), items...)
items = items[:0]
return flush(batch)
}
for {
select {
case <-ctx.Done():
return errors.Join(flushItems(), ctx.Err())
case item, ok := <-input:
if !ok {
return flushItems()
}
items = append(items, item)
if len(items) >= maxSize {
if err := flushItems(); err != nil {
return err
}
}
case <-ticker.C:
if err := flushItems(); err != nil {
return err
}
}
}
}
批量越大,吞吐通常越高,但单条任务等待时间也越长,需要在吞吐和延迟之间权衡。
避免共享内存竞态
以下计数器存在数据竞争:
count := 0
for range jobs {
go func() {
count++
}()
}
可以使用互斥锁、原子操作,或者把状态收敛到单个 Goroutine。不要仅凭“测试结果看起来正确”判断并发安全。
运行竞态检测:
go test -race ./...
go run -race ./cmd/server
-race 会增加运行开销,适合测试和预发布环境,不建议长期用于生产流量。
优雅退出后台任务
应用需要停止接收新任务、取消后台工作,并等待已有任务在限定时间内结束。
func main() {
ctx, stop := signal.NotifyContext(
context.Background(),
os.Interrupt,
syscall.SIGTERM,
)
defer stop()
group, ctx := errgroup.WithContext(ctx)
group.Go(func() error {
return consume(ctx)
})
group.Go(func() error {
return reportMetrics(ctx)
})
if err := group.Wait(); err != nil && !errors.Is(err, context.Canceled) {
log.Printf("service stopped: %v", err)
}
}
对于必须完成的任务,应先停止入口,再等待队列排空。对于可重试任务,应在消息确认前确保处理成功,让未完成任务在重启后重新投递。
监控高并发系统
至少记录以下指标:
- 当前 Goroutine 数量。
- 活跃 Worker 数和任务执行时间。
- 队列长度、容量和入队等待时间。
- 接受、拒绝、失败和重试的任务数。
- 下游连接池使用率与等待时间。
- 限流命中数和超时数。
- 服务关闭阶段剩余任务数。
仅观察平均延迟会掩盖尾部问题,应同时关注 P95、P99 和最大值。
压测与容量评估
压测应逐级增加并发,并同时观察吞吐、延迟、错误率和资源使用:
go test -run '^$' -bench . -benchmem ./...
go test -race ./...
HTTP 压测可以使用 vegeta、wrk 或 k6。压测环境应包含真实的连接池限制和接近生产的数据规模,否则结果只反映 Handler 本身。
找到吞吐不再增长、延迟开始陡升的拐点,并把生产限额设置在拐点之前,而不是把服务推到资源完全耗尽。
实践检查清单
- 并发任务数和队列容量是否都有明确上限。
- 队列满时是否定义了等待、拒绝、丢弃或持久化策略。
- 每个外部调用是否带有超时和取消信号。
- 限流维度是否与用户和业务边界一致。
- 相同键的并发加载是否会造成缓存击穿。
- 共享状态是否通过锁、原子操作或单一所有者保护。
- 是否定期运行 Race Detector 和压力测试。
- 进程退出时是否停止入口、取消任务并等待清理。
- 是否监控队列、Worker、下游连接池和尾部延迟。
高并发系统的核心不是创建更多 Goroutine,而是让流量在每一层都可控。当并发、队列、超时和下游容量形成闭环后,系统才能在正常流量下高效运行,在突发流量下有序降级。