返回首页

Golang进阶:2.3 并发模式与最佳实践

并发是 Go 最具代表性的能力之一。很多开发者刚接触 Go 时,会被 goroutine 和 channel 的简洁写法吸引;但真正进入业务开发后,大家很快会发现:会写并发,不等于写得稳、写得准、写得可维护

这一章我们聚焦几个最常用、也最容易在工程中踩坑的并发主题:

  • context 包:超时控制、取消传播、值传递
  • sync 包:MutexRWMutexOncePoolMap
  • Worker Pool 模式
  • 限流器实现(Rate Limiter)
  • 优雅退出(Graceful Shutdown)

本文不会停留在 API 罗列层面,而是尽量从“为什么这样设计”和“实际项目里怎么用”两个角度展开。默认你已经具备基础的 Go 语法和 goroutine 使用经验。

一、为什么 Go 并发容易写,也容易写错

Go 的并发模型强调“不要通过共享内存来通信,而要通过通信来共享内存”。这句话很经典,但在真实项目里,事情往往没有这么单纯。

例如下面这些场景都很常见:

  • 一个请求需要调用多个下游服务,任何一个超时都要整体取消
  • 一个缓存需要高并发读、少量写
  • 一个对象初始化成本很高,但只能初始化一次
  • 一个后台任务系统要限制同时处理的任务数
  • 一个 HTTP 服务收到退出信号后,要停止接收新请求,并等待已有任务完成

如果没有统一的并发控制策略,代码很容易出现以下问题:

问题 常见表现 风险
goroutine 泄漏 子任务无人回收、阻塞在 channel 或 IO 上 长期运行后内存与 goroutine 数持续上涨
锁使用不当 锁粒度过大、忘记解锁、读写竞争 吞吐下降,甚至死锁
无边界并发 每个请求都直接起 goroutine 峰值时拖垮 CPU、内存或下游系统
缺乏退出机制 程序退出时任务强制中断 数据丢失、状态不一致
上下文丢失 超时、trace、用户信息未向下传递 排查困难、链路不完整

所以,并发的关键不只是“能跑起来”,而是:

  1. 边界清晰:知道任务何时开始、何时停止。
  2. 资源可控:知道最多起多少 goroutine、占用多少资源。
  3. 失败可传播:一个环节失败,相关任务能及时取消。
  4. 退出可收敛:程序停止时,能让已有任务安全完成或尽快收尾。

下面我们从最基础、也最重要的 context 开始。

二、context 包:超时控制、取消传播、值传递

context 是 Go 并发编程中的“控制平面”。它本身不负责执行业务逻辑,但它负责告诉你的 goroutine:

  • 还能不能继续做
  • 最晚可以做到什么时候
  • 是否上游已经取消
  • 当前调用链需要携带哪些请求范围内的数据

2.1 context 的核心作用

标准库里的 context.Context 接口很小,但非常关键:

  • Deadline():返回截止时间
  • Done():返回只读 channel,表示取消信号
  • Err():返回取消原因
  • Value():读取上下文值

在工程实践里,context 最常见的三个用途就是:

用途 说明 典型场景
超时控制 给操作设置最晚完成时间 DB 查询、HTTP 调用、RPC 请求
取消传播 父任务取消时,子任务同步停止 批量并发请求、流水线任务
值传递 在调用链中传递请求级数据 traceID、requestID、用户身份

需要注意:context 适合传递请求范围内的数据,不适合传业务参数,更不适合当作“万能 map”使用。

2.2 超时控制:避免任务无限等待

当一个操作依赖外部资源时,例如数据库、缓存、第三方 API,如果不设置超时,就可能把 goroutine 永久挂住。

下面是一个完整可运行示例,演示如何使用 context.WithTimeout 为任务设置超时。

package main

import (
	"context"
	"fmt"
	"time"
)

func queryRemoteService(ctx context.Context) (string, error) {
	select {
	case <-time.After(3 * time.Second):
		return "service response", nil
	case <-ctx.Done():
		return "", ctx.Err()
	}
}

func main() {
	ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
	defer cancel()

	result, err := queryRemoteService(ctx)
	if err != nil {
		fmt.Println("query failed:", err)
		return
	}

	fmt.Println("query success:", result)
}

这段程序会输出超时错误,因为任务需要 3 秒完成,而上下文只允许 2 秒。

这里有两个实践要点:

  1. 谁创建,谁负责 cancel:即使超时会自动触发,也建议显式 defer cancel(),及时释放关联资源。
  2. 被调用方必须监听 ctx.Done():只有创建超时还不够,下游函数也要配合响应取消。

2.3 取消传播:父任务取消,子任务同步停止

在一个请求里,往往会并发启动多个子任务。例如同时查询用户信息、订单信息、优惠信息。一旦用户主动取消请求,或主流程已经确定失败,就应该尽快通知所有子任务停止。

下面这个例子展示了 context.WithCancel 的取消传播机制。

package main

import (
	"context"
	"fmt"
	"sync"
	"time"
)

func worker(ctx context.Context, id int, wg *sync.WaitGroup) {
	defer wg.Done()

	for {
		select {
		case <-ctx.Done():
			fmt.Printf("worker %d stopped: %v\n", id, ctx.Err())
			return
		default:
			fmt.Printf("worker %d is working\n", id)
			time.Sleep(500 * time.Millisecond)
		}
	}
}

func main() {
	ctx, cancel := context.WithCancel(context.Background())
	var wg sync.WaitGroup

	for i := 1; i <= 3; i++ {
		wg.Add(1)
		go worker(ctx, i, &wg)
	}

	time.Sleep(2 * time.Second)
	fmt.Println("main: cancel all workers")
	cancel()

	wg.Wait()
	fmt.Println("all workers exited")
}

执行后你会看到,所有 worker 都在收到取消信号后退出。这样的设计在服务端程序里很重要,因为它可以避免子 goroutine 在主流程结束后继续“失控运行”。

2.4 值传递:在调用链中携带请求级信息

context.WithValue 常用于跨函数、跨组件传递一些与当前请求强相关的信息,比如:

  • trace ID
  • request ID
  • 登录用户 ID
  • 灰度标记
  • 多租户信息

完整示例如下:

package main

import (
	"context"
	"fmt"
)

type contextKey string

const requestIDKey contextKey = "requestID"

func handler(ctx context.Context) {
	service(ctx)
}

func service(ctx context.Context) {
	repo(ctx)
}

func repo(ctx context.Context) {
	requestID, ok := ctx.Value(requestIDKey).(string)
	if !ok {
		fmt.Println("requestID not found")
		return
	}

	fmt.Println("current requestID:", requestID)
}

func main() {
	ctx := context.WithValue(context.Background(), requestIDKey, "req-20260608-001")
	handler(ctx)
}

这里还体现了一个常见最佳实践:自定义 key 类型。不要直接使用字符串作为 key,否则不同包之间可能发生命名冲突。

2.5 context 的常见误区

误区 问题说明 正确做法
context 存到结构体里长期复用 生命周期不清晰,容易误传取消状态 优先按函数参数逐层传递
忘记调用 cancel 定时器和关联资源不能及时释放 创建后尽量 defer cancel()
context.Value 传业务参数 可读性差,缺乏类型约束 业务参数用显式函数参数传递
下游不监听 ctx.Done() 上游取消后,任务仍继续运行 在阻塞点或循环中检查取消信号
所有函数都强行塞 context 纯计算逻辑并不总需要它 只有请求链、IO链、任务控制链再传

一句话总结:context 不是业务数据容器,而是并发任务的生命周期控制器。

三、sync 包:Mutex、RWMutex、Once、Pool、Map

Go 提倡通过 channel 通信,但并不意味着锁不重要。只要你有共享状态,就很可能会用到 sync 包。

3.1 sync.Mutex:最基础也最常用的互斥锁

Mutex 用来保护共享资源,保证同一时刻只有一个 goroutine 能进入临界区。

下面的示例演示如何安全地对共享计数器累加。

package main

import (
	"fmt"
	"sync"
)

type Counter struct {
	mu    sync.Mutex
	value int
}

func (c *Counter) Inc() {
	c.mu.Lock()
	defer c.mu.Unlock()
	c.value++
}

func (c *Counter) Value() int {
	c.mu.Lock()
	defer c.mu.Unlock()
	return c.value
}

func main() {
	var wg sync.WaitGroup
	counter := &Counter{}

	for i := 0; i < 1000; i++ {
		wg.Add(1)
		go func() {
			defer wg.Done()
			counter.Inc()
		}()
	}

	wg.Wait()
	fmt.Println("final counter:", counter.Value())
}

如果去掉锁,这段代码在高并发下很容易出现竞态,最终结果通常小于 1000。

Mutex 使用建议

  • 临界区尽量小,不要把耗时操作放在锁里
  • 加锁和解锁尽量成对出现,推荐 defer Unlock() 保证异常路径不漏解锁
  • 不要在持锁期间调用外部未知逻辑,避免死锁和性能问题

3.2 sync.RWMutex:读多写少场景的优化选择

RWMutex 区分读锁和写锁:

  • 读锁之间可以并发
  • 写锁是互斥的
  • 写锁会阻塞读锁

适合“读多写少”的场景,例如本地配置缓存、热点字典、规则表等。

package main

import (
	"fmt"
	"sync"
	"time"
)

type ConfigStore struct {
	mu   sync.RWMutex
	data map[string]string
}

func NewConfigStore() *ConfigStore {
	return &ConfigStore{
		data: map[string]string{
			"app_name": "go-blog",
			"env":      "prod",
		},
	}
}

func (s *ConfigStore) Get(key string) (string, bool) {
	s.mu.RLock()
	defer s.mu.RUnlock()
	value, ok := s.data[key]
	return value, ok
}

func (s *ConfigStore) Set(key, value string) {
	s.mu.Lock()
	defer s.mu.Unlock()
	s.data[key] = value
}

func main() {
	store := NewConfigStore()
	var wg sync.WaitGroup

	for i := 0; i < 5; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			for j := 0; j < 3; j++ {
				if value, ok := store.Get("app_name"); ok {
					fmt.Printf("reader %d got app_name=%s\n", id, value)
				}
				time.Sleep(100 * time.Millisecond)
			}
		}(i)
	}

	wg.Add(1)
	go func() {
		defer wg.Done()
		time.Sleep(150 * time.Millisecond)
		store.Set("app_name", "go-blog-v2")
		fmt.Println("writer updated app_name")
	}()

	wg.Wait()
}

不过也要注意,RWMutex 并不是总比 Mutex 更快。对于竞争不高、逻辑简单的场景,Mutex 反而可能更直接、更高效。先按业务特点选,再按性能数据验证。

3.3 sync.Once:只执行一次

Once 常用于全局资源初始化,例如:

  • 单例对象创建
  • 配置加载
  • 连接池初始化
  • 懒加载缓存
package main

import (
	"fmt"
	"sync"
)

type Config struct {
	AppName string
}

var (
	config *Config
	once   sync.Once
)

func GetConfig() *Config {
	once.Do(func() {
		fmt.Println("loading config...")
		config = &Config{AppName: "go-advanced-blog"}
	})
	return config
}

func main() {
	var wg sync.WaitGroup

	for i := 0; i < 5; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			cfg := GetConfig()
			fmt.Printf("goroutine %d got config: %+v\n", id, cfg)
		}(i)
	}

	wg.Wait()
}

无论多少 goroutine 并发调用,初始化逻辑都只会执行一次。

3.4 sync.Pool:复用临时对象,降低分配压力

sync.Pool 适合缓存“可重复使用的临时对象”,比如:

  • bytes.Buffer
  • 临时切片
  • 编解码对象
  • 序列化中间结构

下面示例演示如何复用 bytes.Buffer

package main

import (
	"bytes"
	"fmt"
	"sync"
)

var bufferPool = sync.Pool{
	New: func() any {
		fmt.Println("allocate new buffer")
		return new(bytes.Buffer)
	},
}

func formatMessage(name string, age int) string {
	buf := bufferPool.Get().(*bytes.Buffer)
	buf.Reset()
	defer bufferPool.Put(buf)

	fmt.Fprintf(buf, "name=%s, age=%d", name, age)
	return buf.String()
}

func main() {
	for i := 0; i < 5; i++ {
		msg := formatMessage("alice", 20+i)
		fmt.Println(msg)
	}
}

sync.Pool 的注意事项

  • 它适合“临时对象”,不适合存放需要长期持有的状态对象
  • 池中的对象可能被 GC 清空,所以不要依赖其“永久缓存”语义
  • 放回池前务必重置状态,例如 Reset(),避免脏数据污染下一次使用

3.5 sync.Map:并发安全的 map

Go 原生 map 不是并发安全的,多 goroutine 读写会直接 panic。sync.Map 提供了并发安全的 map 实现,适合以下场景:

  • key 集合相对稳定,读远多于写
  • 多 goroutine 独立读写,冲突模式比较复杂
  • 希望省掉手动加锁的模板代码
package main

import (
	"fmt"
	"sync"
)

func main() {
	var m sync.Map
	var wg sync.WaitGroup

	for i := 0; i < 5; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			key := fmt.Sprintf("user:%d", id)
			m.Store(key, id*10)
		} (i)
	}

	wg.Wait()

	m.Range(func(key, value any) bool {
		fmt.Printf("%v => %v\n", key, value)
		return true
	})

	if value, ok := m.Load("user:3"); ok {
		fmt.Println("user:3 =", value)
	}
}

如果你的场景需要复杂事务逻辑、批量操作、强类型封装,那么“map + mutex”通常更容易维护。sync.Map 不是普通 map 的无脑替代品。

3.6 如何选择合适的 sync 原语

工具 适用场景 优点 注意点
Mutex 通用共享状态保护 简单直接 锁粒度要合理
RWMutex 读多写少 提高读并发 不一定总比 Mutex
Once 单次初始化 天然并发安全 初始化失败需自行设计补偿策略
Pool 临时对象复用 降低分配和 GC 压力 对象放回前要清理状态
Map 并发 map 访问 免手动加锁 不适合所有数据结构场景

四、Worker Pool 模式

在很多业务系统中,任务的产生速度通常快于处理速度。例如:

  • 消费消息队列中的任务
  • 批量处理用户导入数据
  • 并发抓取网页或调用第三方接口
  • 图片、音视频、日志等异步处理任务

如果每来一个任务就启动一个 goroutine,短时间内虽然很“爽”,但一旦任务量暴涨,就可能出现:

  • goroutine 数量失控
  • CPU 调度压力过大
  • 内存占用持续上涨
  • 下游系统被并发打爆

Worker Pool 的核心思路是:任务可以很多,但同时工作的 worker 数量必须可控。

4.1 Worker Pool 的基本结构

一个典型的 Worker Pool 一般包含:

  1. 任务队列
  2. 固定数量的 worker
  3. 结果通道(可选)
  4. 退出控制机制

下面先看一个结构清晰、接近生产使用方式的完整示例。它支持:

  • 固定 worker 数量
  • 提交任务
  • 收集结果
  • 基于 context 的取消控制
  • 优雅关闭任务队列并等待 worker 退出
package main

import (
	"context"
	"fmt"
	"math/rand"
	"sync"
	"time"
)

type Job struct {
	ID      int
	Payload int
}

type Result struct {
	JobID int
	Value int
	Err   error
}

type Pool struct {
	jobs    chan Job
	results chan Result
	wg      sync.WaitGroup
}

func NewPool(workerCount, queueSize int) *Pool {
	p := &Pool{
		jobs:    make(chan Job, queueSize),
		results: make(chan Result, queueSize),
	}

	for i := 1; i <= workerCount; i++ {
		p.wg.Add(1)
		go p.worker(i)
	}

	return p
}

func (p *Pool) worker(workerID int) {
	defer p.wg.Done()

	for job := range p.jobs {
		fmt.Printf("worker %d processing job %d\n", workerID, job.ID)
		time.Sleep(time.Duration(rand.Intn(500)+300) * time.Millisecond)

		result := Result{
			JobID: job.ID,
			Value: job.Payload * 2,
		}
		p.results <- result
	}

	fmt.Printf("worker %d exited\n", workerID)
}

func (p *Pool) Submit(ctx context.Context, job Job) error {
	select {
	case <-ctx.Done():
		return ctx.Err()
	case p.jobs <- job:
		return nil
	}
}

func (p *Pool) Results() <-chan Result {
	return p.results
}

func (p *Pool) Close() {
	close(p.jobs)
	p.wg.Wait()
	close(p.results)
}

func main() {
	rand.Seed(time.Now().UnixNano())

	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel()

	pool := NewPool(3, 10)

	go func() {
		for i := 1; i <= 8; i++ {
			job := Job{ID: i, Payload: i * 10}
			if err := pool.Submit(ctx, job); err != nil {
				fmt.Printf("submit job %d failed: %v\n", i, err)
				continue
			}
		}
		pool.Close()
	}()

	for result := range pool.Results() {
		if result.Err != nil {
			fmt.Printf("job %d failed: %v\n", result.JobID, result.Err)
			continue
		}
		fmt.Printf("job %d result: %d\n", result.JobID, result.Value)
	}

	fmt.Println("all jobs done")
}

4.2 生产环境下 Worker Pool 需要关注什么

上面这个例子已经具备实用雏形,但真实生产中还要考虑更多因素:

关注点 说明
队列长度 防止任务无限堆积,必要时拒绝新任务
超时控制 单个任务不能无限执行,最好配合 context
错误处理 区分可重试错误和不可重试错误
指标监控 监控活跃 worker 数、队列长度、失败率、处理耗时
背压策略 当生产速度超过消费能力时,要能限流或丢弃
关闭流程 先停止接收新任务,再等待已有任务处理完成

4.3 更贴近生产环境的 Worker Pool 示例

下面这个版本更接近服务端项目中的写法。它支持:

  • worker 生命周期绑定到 context
  • 任务处理超时
  • 非阻塞结果上报
  • 停止接收新任务
  • 安全关闭
package main

import (
	"context"
	"errors"
	"fmt"
	"math/rand"
	"sync"
	"sync/atomic"
	"time"
)

type Task struct {
	ID   int
	Name string
}

type TaskResult struct {
	TaskID   int
	WorkerID int
	Err      error
	Cost     time.Duration
}

type WorkerPool struct {
	ctx      context.Context
	cancel   context.CancelFunc
	jobs     chan Task
	results  chan TaskResult
	wg       sync.WaitGroup
	closed   atomic.Bool
	once     sync.Once
	workers  int
	timeout  time.Duration
}

func NewWorkerPool(parent context.Context, workerCount, queueSize int, taskTimeout time.Duration) *WorkerPool {
	ctx, cancel := context.WithCancel(parent)
	p := &WorkerPool{
		ctx:     ctx,
		cancel:  cancel,
		jobs:    make(chan Task, queueSize),
		results: make(chan TaskResult, queueSize),
		workers: workerCount,
		timeout: taskTimeout,
	}

	for i := 1; i <= workerCount; i++ {
		p.wg.Add(1)
		go p.runWorker(i)
	}

	return p
}

func (p *WorkerPool) runWorker(workerID int) {
	defer p.wg.Done()

	for {
		select {
		case <-p.ctx.Done():
			fmt.Printf("worker %d stopped: %v\n", workerID, p.ctx.Err())
			return
		case task, ok := <-p.jobs:
			if !ok {
				fmt.Printf("worker %d exit: job channel closed\n", workerID)
				return
			}

			start := time.Now()
			err := p.handleTask(task)
			result := TaskResult{
				TaskID:   task.ID,
				WorkerID: workerID,
				Err:      err,
				Cost:     time.Since(start),
			}

			select {
			case p.results <- result:
			case <-p.ctx.Done():
				return
			}
		}
	}
}

func (p *WorkerPool) handleTask(task Task) error {
	taskCtx, cancel := context.WithTimeout(p.ctx, p.timeout)
	defer cancel()

	workTime := time.Duration(rand.Intn(1500)+300) * time.Millisecond

	select {
	case <-time.After(workTime):
		if workTime > 1200*time.Millisecond {
			return errors.New("task execution too slow, marked as failed")
		}
		return nil
	case <-taskCtx.Done():
		return taskCtx.Err()
	}
}

func (p *WorkerPool) Submit(task Task) error {
	if p.closed.Load() {
		return errors.New("worker pool is closed")
	}

	select {
	case <-p.ctx.Done():
		return p.ctx.Err()
	case p.jobs <- task:
		return nil
	default:
		return errors.New("job queue is full")
	}
}

func (p *WorkerPool) Results() <-chan TaskResult {
	return p.results
}

func (p *WorkerPool) Shutdown() {
	p.once.Do(func() {
		p.closed.Store(true)
		close(p.jobs)
		p.wg.Wait()
		close(p.results)
		p.cancel()
	})
}

func main() {
	rand.Seed(time.Now().UnixNano())

	ctx := context.Background()
	pool := NewWorkerPool(ctx, 4, 8, 1*time.Second)

	for i := 1; i <= 12; i++ {
		err := pool.Submit(Task{ID: i, Name: fmt.Sprintf("task-%d", i)})
		if err != nil {
			fmt.Printf("submit failed for task %d: %v\n", i, err)
		}
	}

	go func() {
		pool.Shutdown()
	}()

	for result := range pool.Results() {
		if result.Err != nil {
			fmt.Printf("task %d handled by worker %d failed, cost=%v, err=%v\n",
				result.TaskID, result.WorkerID, result.Cost, result.Err)
			continue
		}

		fmt.Printf("task %d handled by worker %d success, cost=%v\n",
			result.TaskID, result.WorkerID, result.Cost)
	}

	fmt.Println("worker pool shutdown complete")
}

这个版本体现了几个很关键的工程思路:

  • 并发不是无限放大,而是被池大小约束
  • 提交失败也是正常业务结果,例如队列满了就应该拒绝
  • 任务处理必须有超时,防止单任务卡死 worker
  • 关闭流程必须有顺序:先关入口,再等 worker 退出,再关结果通道

五、限流器实现(Rate Limiter)

在并发系统中,仅仅控制“同时跑多少任务”还不够,很多时候还要控制“单位时间内允许多少请求”。这就是限流。

限流常见于以下场景:

  • 接口防刷
  • 保护数据库和缓存
  • 控制第三方 API 调用频率
  • 下游服务容量保护
  • 消息消费速率整形

5.1 常见限流算法

算法 特点 适用场景
固定窗口 实现简单 粗粒度限流
滑动窗口 更平滑 网关、接口限流
漏桶 输出速率稳定 平滑流量整形
令牌桶 允许一定突发流量 服务端最常见

在 Go 项目里,令牌桶是非常常见的选择。它的核心思想是:

  • 桶中持续按固定速率产生令牌
  • 请求到来时先尝试拿令牌
  • 拿到就通过,拿不到就拒绝或等待
  • 桶容量决定了最大突发能力

5.2 基于 channel 的简单令牌桶实现

下面先看一个容易理解、完整可运行的版本。

package main

import (
	"fmt"
	"time"
)

type RateLimiter struct {
	tokens chan struct{}
}

func NewRateLimiter(rate int, burst int) *RateLimiter {
	rl := &RateLimiter{
		tokens: make(chan struct{}, burst),
	}

	for i := 0; i < burst; i++ {
		rl.tokens <- struct{}{}
	}

	interval := time.Second / time.Duration(rate)
	ticker := time.NewTicker(interval)

	go func() {
		for range ticker.C {
			select {
			case rl.tokens <- struct{}{}:
			default:
				// 桶满了,跳过
			}
		}
	}()

	return rl
}

func (rl *RateLimiter) Allow() bool {
	select {
	case <-rl.tokens:
		return true
	default:
		return false
	}
}

func main() {
	rl := NewRateLimiter(2, 4) // 每秒 2 个令牌,桶容量 4

	for i := 1; i <= 10; i++ {
		if rl.Allow() {
			fmt.Printf("request %d allowed at %s\n", i, time.Now().Format("15:04:05"))
		} else {
			fmt.Printf("request %d rejected at %s\n", i, time.Now().Format("15:04:05"))
		}
		time.Sleep(300 * time.Millisecond)
	}
}

这个实现很适合教学和小型场景。它的优点是直观,缺点是功能较基础,比如:

  • 不支持等待直到可用
  • 不支持和 context 联动取消等待
  • 不方便扩展多维限流

5.3 支持等待与超时的限流器实现

更实用的限流器往往不只是“允许/拒绝”,还需要支持:

  • 等待令牌
  • 超时取消
  • 优雅停止

下面给出一个更完整的实现示例。

package main

import (
	"context"
	"fmt"
	"sync"
	"time"
)

type TokenBucket struct {
	tokens chan struct{}
	stopCh chan struct{}
	once   sync.Once
}

func NewTokenBucket(rate int, burst int) *TokenBucket {
	tb := &TokenBucket{
		tokens: make(chan struct{}, burst),
		stopCh: make(chan struct{}),
	}

	for i := 0; i < burst; i++ {
		tb.tokens <- struct{}{}
	}

	interval := time.Second / time.Duration(rate)
	ticker := time.NewTicker(interval)

	go func() {
		defer ticker.Stop()
		for {
			select {
			case <-tb.stopCh:
				return
			case <-ticker.C:
				select {
				case tb.tokens <- struct{}{}:
				default:
				}
			}
		}
	}()

	return tb
}

func (tb *TokenBucket) Allow() bool {
	select {
	case <-tb.tokens:
		return true
	default:
		return false
	}
}

func (tb *TokenBucket) Wait(ctx context.Context) error {
	select {
	case <-ctx.Done():
		return ctx.Err()
	case <-tb.tokens:
		return nil
	}
}

func (tb *TokenBucket) Stop() {
	tb.once.Do(func() {
		close(tb.stopCh)
	})
}

func main() {
	limiter := NewTokenBucket(3, 3) // 每秒 3 个令牌,最大突发 3
	defer limiter.Stop()

	for i := 1; i <= 8; i++ {
		ctx, cancel := context.WithTimeout(context.Background(), 700*time.Millisecond)
		err := limiter.Wait(ctx)
		cancel()

		if err != nil {
			fmt.Printf("request %d timeout: %v\n", i, err)
			continue
		}

		fmt.Printf("request %d passed at %s\n", i, time.Now().Format("15:04:05.000"))
	}
}

这个版本已经具备较好的工程可用性。它特别适合挂在接口调用、异步任务分发、第三方 API 客户端前面。

5.4 限流与并发控制的区别

很多初学者会把“限流”和“并发控制”混为一谈,但两者控制的是不同维度:

维度 控制目标 典型手段
并发控制 同一时刻最多多少任务同时执行 Worker Pool、信号量
速率控制 单位时间最多允许多少次请求 Rate Limiter

实际项目中,这两者往往需要结合使用。例如:

  • 用 Worker Pool 控制最多只有 20 个任务同时执行
  • 用 Rate Limiter 控制每秒最多只能调用下游接口 100 次

这样才能同时保护本机资源和下游资源。

六、优雅退出(Graceful Shutdown)

优雅退出,是并发程序走向工程化的关键一步。

很多程序在本地测试时,习惯直接按 Ctrl+C 结束进程,看起来很正常;但在生产环境里,如果没有优雅退出机制,可能导致:

  • 正在处理的请求被强行中断
  • 任务只执行到一半,状态不一致
  • 数据尚未刷盘或提交
  • 日志、监控、trace 来不及上报
  • 数据库连接、消息消费、定时任务没有正常关闭

所谓优雅退出,本质上就是:

  1. 收到退出信号
  2. 停止接收新流量或新任务
  3. 通知已有 goroutine 进入收尾阶段
  4. 等待它们在限定时间内完成
  5. 超时后再强制退出

6.1 优雅退出的核心步骤

步骤 说明
捕获信号 监听 SIGINTSIGTERM
广播取消 通过 context 或关闭 channel 通知子任务退出
停止入口 停止 HTTP 服务、消息消费、任务投递
等待收尾 等待正在执行的任务完成
超时兜底 超时后放弃等待,避免卡死退出

6.2 一个基础可运行示例

先看一个基础版本,演示如何监听信号并让后台任务退出。

package main

import (
	"context"
	"fmt"
	"os"
	"os/signal"
	"sync"
	"syscall"
	"time"
)

func backgroundWorker(ctx context.Context, wg *sync.WaitGroup, name string) {
	defer wg.Done()
	ticker := time.NewTicker(1 * time.Second)
	defer ticker.Stop()

	for {
		select {
		case <-ctx.Done():
			fmt.Println(name, "received stop signal")
			return
		case t := <-ticker.C:
			fmt.Println(name, "running at", t.Format("15:04:05"))
		}
	}
}

func main() {
	ctx, cancel := context.WithCancel(context.Background())
	var wg sync.WaitGroup

	wg.Add(2)
	go backgroundWorker(ctx, &wg, "worker-1")
	go backgroundWorker(ctx, &wg, "worker-2")

	sigCh := make(chan os.Signal, 1)
	signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)

	sig := <-sigCh
	fmt.Println("received signal:", sig)
	cancel()

	wg.Wait()
	fmt.Println("graceful shutdown complete")
}

这个版本适合理解基本流程,但它还不够“生产环境化”,因为它没有:

  • HTTP 服务关闭流程
  • 超时控制
  • 在途请求等待
  • 后台任务统一收口

6.3 接近生产环境的 Graceful Shutdown 完整示例

下面我们实现一个更完整的示例:

  • 启动一个 HTTP 服务
  • 每个请求模拟业务处理
  • 后台还有一个定时任务协程
  • 收到退出信号后,停止接收新请求
  • 给已有请求和后台任务一个收尾时间窗口
  • 使用 context 统一传递取消信号
package main

import (
	"context"
	"fmt"
	"log"
	"net/http"
	"os"
	"os/signal"
	"sync"
	"syscall"
	"time"
)

type App struct {
	server *http.Server
	wg     sync.WaitGroup
}

func NewApp(addr string) *App {
	app := &App{}
	mux := http.NewServeMux()

	mux.HandleFunc("/work", app.handleWork)
	mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
		w.WriteHeader(http.StatusOK)
		_, _ = w.Write([]byte("ok"))
	})

	app.server = &http.Server{
		Addr:    addr,
		Handler: mux,
	}

	return app
}

func (a *App) handleWork(w http.ResponseWriter, r *http.Request) {
	ctx := r.Context()
	log.Println("request started")

	select {
	case <-time.After(3 * time.Second):
		w.WriteHeader(http.StatusOK)
		_, _ = w.Write([]byte("work done"))
		log.Println("request finished")
	case <-ctx.Done():
		log.Println("request canceled:", ctx.Err())
		http.Error(w, "request canceled", http.StatusRequestTimeout)
	}
}

func (a *App) startBackgroundJob(ctx context.Context) {
	a.wg.Add(1)
	go func() {
		defer a.wg.Done()
		ticker := time.NewTicker(2 * time.Second)
		defer ticker.Stop()

		for {
			select {
			case <-ctx.Done():
				log.Println("background job stopping...")
				// 模拟清理动作
				time.Sleep(1 * time.Second)
				log.Println("background job stopped")
				return
			case t := <-ticker.C:
				log.Println("background job tick at", t.Format("15:04:05"))
			}
		}
	}()
}

func (a *App) Run() error {
	rootCtx, cancel := context.WithCancel(context.Background())
	defer cancel()

	a.startBackgroundJob(rootCtx)

	serverErrCh := make(chan error, 1)
	go func() {
		log.Println("http server listening on", a.server.Addr)
		if err := a.server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
			serverErrCh <- err
		}
	}()

	sigCh := make(chan os.Signal, 1)
	signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)

	select {
	case err := <-serverErrCh:
		return fmt.Errorf("server error: %w", err)
	case sig := <-sigCh:
		log.Println("received signal:", sig)
	}

	cancel()

	shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer shutdownCancel()

	log.Println("shutting down http server...")
	if err := a.server.Shutdown(shutdownCtx); err != nil {
		return fmt.Errorf("http shutdown failed: %w", err)
	}

	waitCh := make(chan struct{})
	go func() {
		defer close(waitCh)
		a.wg.Wait()
	}()

	select {
	case <-waitCh:
		log.Println("all background jobs stopped")
	case <-shutdownCtx.Done():
		return fmt.Errorf("background jobs shutdown timeout: %w", shutdownCtx.Err())
	}

	log.Println("graceful shutdown complete")
	return nil
}

func main() {
	app := NewApp(":8080")
	if err := app.Run(); err != nil {
		log.Fatal(err)
	}
}

这段代码已经很接近真实服务程序的退出逻辑了。你可以运行后访问 http://localhost:8080/work,然后在请求处理中发送退出信号,观察服务如何进入收尾状态。

6.4 Graceful Shutdown 的实践建议

建议 说明
统一根 context 所有后台任务都挂到同一个根上下文之下
先停入口,再等存量 先停止接收新请求,再等待已有请求完成
设置退出超时 防止某个任务卡死导致整个进程无法退出
ShutdownClose 区分清楚 Shutdown 尽量温和,Close 更偏强制
后台任务要可取消 定时任务、消费者、worker 都必须能响应退出信号
打印关键日志 方便排查退出时卡在哪一步

七、并发最佳实践总结

学完上面的工具和模式后,还需要建立一套“并发设计习惯”。下面这些经验在 Go 项目里非常重要。

7.1 优先考虑生命周期管理

很多并发 bug 的根源不是“语法不会”,而是“生命周期没设计清楚”。

在启动一个 goroutine 之前,最好先回答这几个问题:

  • 它由谁启动?
  • 它何时结束?
  • 上游取消时它会不会停?
  • 它失败了要不要上报?
  • 它退出时是否需要清理资源?

只要这些问题答不上来,代码后面大概率会出问题。

7.2 不要无限制地启动 goroutine

goroutine 很轻量,但绝不是免费。真正的工程原则是:

并发能力必须有边界。

常见手段包括:

  • Worker Pool
  • 带缓冲 channel
  • 信号量
  • Rate Limiter
  • 有界队列

7.3 用 context 管控制,用参数管业务

一个很实用的判断标准是:

  • 与取消、超时、链路信息相关的,用 context
  • 与业务处理逻辑相关的,用普通参数或结构体字段

这样代码职责更清晰,也更容易维护。

7.4 锁要为数据服务,而不是为“安全感”服务

不要看到共享变量就条件反射加大锁。要先想清楚:

  • 这个共享状态是否真的必须共享?
  • 能否改成消息传递?
  • 锁的粒度是否过大?
  • 是否读多写少,适合 RWMutex
  • 是否只是临时对象复用,适合 sync.Pool

并发优化很多时候不是“加更多机制”,而是“减少共享”。

7.5 让退出流程成为系统设计的一部分

优雅退出不是上线前补的一层壳,而应该从一开始就纳入设计。特别是下面这些组件:

  • HTTP / RPC 服务
  • 消息消费者
  • 定时任务
  • 后台异步 worker
  • 数据刷盘与资源回收逻辑

如果这些模块都没有退出协议,那么程序越大,退出越容易乱。

八、结语

Go 给了我们非常好用的并发工具:goroutine 足够轻,channel 足够直观,contextsync 也足够实用。但真正写出稳定的并发程序,靠的不是单个语法点,而是整体设计能力。

这一章你需要真正掌握的,不只是几个 API 名字,而是这几种思维方式:

  • context 管任务生命周期
  • sync 安全管理共享状态
  • 用 Worker Pool 控制并发边界
  • 用 Rate Limiter 控制访问速率
  • 用 Graceful Shutdown 让系统平稳收尾

当这些能力组合起来之后,你写出来的 Go 程序才会从“能并发”走向“并发可控、行为可预期、线上可维护”。

下一步如果你愿意继续深入,可以开始系统学习:channel 编排、pipeline 模式、errgroup、并发调试、竞态检测与性能分析。这些内容会进一步帮助你把 Go 的并发能力用到更成熟的工程实践中。


📝 版权声明:本文为原创技术博客,转载请注明出处。

如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!

上一篇

Golang进阶:2.2 Channel 深度解析

下一篇

Golang进阶:2.4 泛型编程(Go 1.18+)