Go 高并发实战:Worker Pool、限流、背压与优雅退出

Go 可以轻松启动 Goroutine,但“能并发”并不等于“能稳定承受高并发”。生产系统还要控制任务数量、内存占用、下游压力、超时传播和停机过程。

本文从一个原则出发:系统中的每个队列和并发点都必须有上限,每个后台任务都必须能退出。

先定义容量而不是先启动 Goroutine

以下写法会为每个任务启动一个 Goroutine:

for _, job := range jobs {
	go handle(job)
}

任务数量较小时没有问题,但面对突发流量可能产生以下后果:

  • Goroutine 数量持续增长,占用栈和调度资源。
  • 所有任务同时请求数据库或下游服务,迅速耗尽连接池。
  • 任务在内存中排队,延迟和内存同时上升。
  • 服务退出时无法确认哪些任务已经完成。

设计并发系统前,应明确四个容量参数:

  1. 同时执行多少任务。
  2. 队列最多容纳多少任务。
  3. 每个任务允许执行多久。
  4. 下游每秒最多接受多少请求。

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 Requests503 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 压测可以使用 vegetawrkk6。压测环境应包含真实的连接池限制和接近生产的数据规模,否则结果只反映 Handler 本身。

找到吞吐不再增长、延迟开始陡升的拐点,并把生产限额设置在拐点之前,而不是把服务推到资源完全耗尽。

实践检查清单

  1. 并发任务数和队列容量是否都有明确上限。
  2. 队列满时是否定义了等待、拒绝、丢弃或持久化策略。
  3. 每个外部调用是否带有超时和取消信号。
  4. 限流维度是否与用户和业务边界一致。
  5. 相同键的并发加载是否会造成缓存击穿。
  6. 共享状态是否通过锁、原子操作或单一所有者保护。
  7. 是否定期运行 Race Detector 和压力测试。
  8. 进程退出时是否停止入口、取消任务并等待清理。
  9. 是否监控队列、Worker、下游连接池和尾部延迟。

高并发系统的核心不是创建更多 Goroutine,而是让流量在每一层都可控。当并发、队列、超时和下游容量形成闭环后,系统才能在正常流量下高效运行,在突发流量下有序降级。