并发是 Go 最具代表性的能力之一。很多开发者刚接触 Go 时,会被 goroutine 和 channel 的简洁写法吸引;但真正进入业务开发后,大家很快会发现:会写并发,不等于写得稳、写得准、写得可维护。
这一章我们聚焦几个最常用、也最容易在工程中踩坑的并发主题:
context包:超时控制、取消传播、值传递sync包:Mutex、RWMutex、Once、Pool、Map- Worker Pool 模式
- 限流器实现(Rate Limiter)
- 优雅退出(Graceful Shutdown)
本文不会停留在 API 罗列层面,而是尽量从“为什么这样设计”和“实际项目里怎么用”两个角度展开。默认你已经具备基础的 Go 语法和 goroutine 使用经验。
一、为什么 Go 并发容易写,也容易写错
Go 的并发模型强调“不要通过共享内存来通信,而要通过通信来共享内存”。这句话很经典,但在真实项目里,事情往往没有这么单纯。
例如下面这些场景都很常见:
- 一个请求需要调用多个下游服务,任何一个超时都要整体取消
- 一个缓存需要高并发读、少量写
- 一个对象初始化成本很高,但只能初始化一次
- 一个后台任务系统要限制同时处理的任务数
- 一个 HTTP 服务收到退出信号后,要停止接收新请求,并等待已有任务完成
如果没有统一的并发控制策略,代码很容易出现以下问题:
| 问题 | 常见表现 | 风险 |
|---|---|---|
| goroutine 泄漏 | 子任务无人回收、阻塞在 channel 或 IO 上 | 长期运行后内存与 goroutine 数持续上涨 |
| 锁使用不当 | 锁粒度过大、忘记解锁、读写竞争 | 吞吐下降,甚至死锁 |
| 无边界并发 | 每个请求都直接起 goroutine | 峰值时拖垮 CPU、内存或下游系统 |
| 缺乏退出机制 | 程序退出时任务强制中断 | 数据丢失、状态不一致 |
| 上下文丢失 | 超时、trace、用户信息未向下传递 | 排查困难、链路不完整 |
所以,并发的关键不只是“能跑起来”,而是:
- 边界清晰:知道任务何时开始、何时停止。
- 资源可控:知道最多起多少 goroutine、占用多少资源。
- 失败可传播:一个环节失败,相关任务能及时取消。
- 退出可收敛:程序停止时,能让已有任务安全完成或尽快收尾。
下面我们从最基础、也最重要的 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 秒。
这里有两个实践要点:
- 谁创建,谁负责 cancel:即使超时会自动触发,也建议显式
defer cancel(),及时释放关联资源。 - 被调用方必须监听
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 一般包含:
- 任务队列
- 固定数量的 worker
- 结果通道(可选)
- 退出控制机制
下面先看一个结构清晰、接近生产使用方式的完整示例。它支持:
- 固定 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 来不及上报
- 数据库连接、消息消费、定时任务没有正常关闭
所谓优雅退出,本质上就是:
- 收到退出信号
- 停止接收新流量或新任务
- 通知已有 goroutine 进入收尾阶段
- 等待它们在限定时间内完成
- 超时后再强制退出
6.1 优雅退出的核心步骤
| 步骤 | 说明 |
|---|---|
| 捕获信号 | 监听 SIGINT、SIGTERM |
| 广播取消 | 通过 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 |
所有后台任务都挂到同一个根上下文之下 |
| 先停入口,再等存量 | 先停止接收新请求,再等待已有请求完成 |
| 设置退出超时 | 防止某个任务卡死导致整个进程无法退出 |
Shutdown 与 Close 区分清楚 |
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 足够直观,context 和 sync 也足够实用。但真正写出稳定的并发程序,靠的不是单个语法点,而是整体设计能力。
这一章你需要真正掌握的,不只是几个 API 名字,而是这几种思维方式:
- 用
context管任务生命周期 - 用
sync安全管理共享状态 - 用 Worker Pool 控制并发边界
- 用 Rate Limiter 控制访问速率
- 用 Graceful Shutdown 让系统平稳收尾
当这些能力组合起来之后,你写出来的 Go 程序才会从“能并发”走向“并发可控、行为可预期、线上可维护”。
下一步如果你愿意继续深入,可以开始系统学习:channel 编排、pipeline 模式、errgroup、并发调试、竞态检测与性能分析。这些内容会进一步帮助你把 Go 的并发能力用到更成熟的工程实践中。
📝 版权声明:本文为原创技术博客,转载请注明出处。
如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!