并发模型
Go 的并发哲学:不要通过共享内存来通信,而要通过通信来共享内存。
核心组件
| 组件 | 作用 | 类比 |
|---|
| goroutine | 轻量级线程 | 绿色线程 / 协程 |
| channel | goroutine 间通信管道 | Unix pipe |
| select | 多路复用 channel 操作 | epoll |
| sync.Mutex | 互斥锁 | 信号量 |
| sync.WaitGroup | 等待一组 goroutine 完成 | CountDownLatch |
| context | 超时 / 取消 / 值传递 | 线程上下文 |
goroutine
// 启动 goroutine:go 关键字
go doSomething()
// 带参数
go func(name string) {
fmt.Println("hello", name)
}("小徐")
// 匿名函数闭包
name := "小徐"
go func() {
fmt.Println("hello", name) // 捕获外部变量
}()
goroutine vs 线程对比
| 对比 | OS 线程 | goroutine |
|---|
| 初始栈大小 | 1-8 MB | 2 KB(可动态伸缩) |
| 创建开销 | ~1ms | ~1μs |
| 调度方式 | 内核抢占式 | Go runtime 用户态调度(GMP 模型) |
| 上下文切换 | ~1μs(内核态切换) | ~100ns(用户态切换) |
| 数量上限 | 数千 | 数十万 |
| 阻塞 | 整个线程阻塞 | runtime 自动切换其他 goroutine |
GMP 调度模型
G (Goroutine) - goroutine 本身
M (Machine) - OS 线程
P (Processor) - 逻辑处理器,持有可运行 G 的本地队列
┌─────────┐
│ Scheduler│
└────┬────┘
│
┌────────┼────────┐
│ │ │
P1 P2 P3 (GOMAXPROCS 个 P)
[G G G] [G G] [G G G G] 本地运行队列
│ │ │
M1 M2 M3 OS 线程
│ │ │
内核 内核 内核
Work Stealing:当 P 的本地队列空了,会从其他 P 偷 G 来执行
channel
创建与操作
// 无缓冲 channel:同步(发送和接收必须同时就绪)
ch := make(chan int)
// 有缓冲 channel:异步(缓冲未满时发送不阻塞)
ch := make(chan int, 10)
// 发送
ch <- 42
// 接收
v := <-ch
// 接收 + 检查是否关闭
v, ok := <-ch
// 关闭
close(ch)
// 遍历(直到 channel 关闭)
for v := range ch {
fmt.Println(v)
}
缓冲 vs 无缓冲
| 特性 | 无缓冲 make(chan T) | 有缓冲 make(chan T, n) |
|---|
| 发送阻塞 | 直到有人接收 | 缓冲满才阻塞 |
| 接收阻塞 | 直到有人发送 | 缓冲空才阻塞 |
| 同步性 | 强同步(握手) | 异步(解耦) |
| 典型用途 | 信号通知、同步等待 | 生产者-消费者、限流 |
channel 状态矩阵
| 操作 | nil channel | 已关闭 channel | 正常 channel |
|---|
发送 ch <- v | 永久阻塞 | panic | 阻塞或成功 |
接收 <-ch | 永久阻塞 | 返回零值(不阻塞) | 阻塞或成功 |
关闭 close(ch) | panic | panic | 成功 |
长度 len(ch) | 0 | 0 | 当前缓冲元素数 |
容量 cap(ch) | 0 | 缓冲容量 | 缓冲容量 |
常见模式
// 模式1:等待 goroutine 完成
func main() {
ch := make(chan struct{})
go func() {
// 执行工作
doWork()
close(ch) // 完成后关闭
}()
<-ch // 阻塞等待完成
}
// 模式2:fan-out(分发任务到多个 worker)
func fanOut(jobs <-chan int, results chan<- int, workerCount int) {
var wg sync.WaitGroup
for i := 0; i < workerCount; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for job := range jobs {
results <- process(job)
}
}()
}
go func() {
wg.Wait()
close(results)
}()
}
// 模式3:fan-in(合并多个 channel)
func fanIn(channels ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
multiplex := func(c <-chan int) {
defer wg.Done()
for v := range c {
out <- v
}
}
for _, c := range channels {
wg.Add(1)
go multiplex(c)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
// 模式4:生产者-消费者
func producerConsumer() {
jobs := make(chan int, 100)
results := make(chan int, 100)
// 生产者
go func() {
for i := 0; i < 1000; i++ {
jobs <- i
}
close(jobs)
}()
// 3 个消费者
for i := 0; i < 3; i++ {
go func(id int) {
for job := range jobs {
results <- job * 2
}
}(i)
}
go func() {
// 等待所有消费者完成后关闭 results
}()
}
// 模式5:限流(带缓冲 channel 作为信号量)
func rateLimit(maxConcurrent int, tasks []func()) {
sem := make(chan struct{}, maxConcurrent)
var wg sync.WaitGroup
for _, task := range tasks {
wg.Add(1)
sem <- struct{}{} // 获取令牌
go func(t func()) {
defer wg.Done()
defer func() { <-sem }() // 释放令牌
t()
}(task)
}
wg.Wait()
}
// 模式6:超时控制
select {
case result := <-ch:
fmt.Println("got result:", result)
case <-time.After(5 * time.Second):
fmt.Println("timeout")
}
// 模式7:心跳检测
func heartbeat(work <-chan int) <-chan struct{} {
hb := make(chan struct{})
go func() {
defer close(hb)
for {
select {
case _, ok := <-work:
if !ok {
return
}
hb <- struct{}{}
case <-time.After(10 * time.Second):
return // 超时退出
}
}
}()
return hb
}
select
// 多路复用:哪个 channel 先就绪就执行哪个
select {
case v := <-ch1:
fmt.Println("from ch1:", v)
case v := <-ch2:
fmt.Println("from ch2:", v)
case ch3 <- 42:
fmt.Println("sent to ch3")
case <-time.After(1 * time.Second):
fmt.Println("timeout")
default:
fmt.Println("no channel ready") // 非阻塞模式
}
| select 行为 | 条件 |
|---|
| 随机选择一个就绪的 case | 多个 case 就绪 |
| 阻塞等待 | 没有 case 就绪且无 default |
| 执行 default | 没有 case 就绪但有 default |
| 立即返回 | 有 case 就绪 |
select 实现超时
// 超时 + 取消
func doWithTimeout(ctx context.Context, timeout time.Duration) (Result, error) {
ctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
resultCh := make(chan Result, 1)
errCh := make(chan error, 1)
go func() {
result, err := slowOperation(ctx)
if err != nil {
errCh <- err
return
}
resultCh <- result
}()
select {
case result := <-resultCh:
return result, nil
case err := <-errCh:
return Result{}, err
case <-ctx.Done():
return Result{}, ctx.Err()
}
}
sync 包核心原语
WaitGroup
var wg sync.WaitGroup
for i := 0; i < 5; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done() // 必须在 goroutine 内 defer
doWork(id)
}(i)
}
wg.Wait() // 等待所有完成
// ⚠️ Add 必须在 goroutine 外调用,Done 必须在 goroutine 内调用
// ⚠️ WaitGroup 不可复制(值传递),必须用指针
Mutex / RWMutex
// 互斥锁
var mu sync.Mutex
mu.Lock()
// 临界区
mu.Unlock()
// 读写锁
var rw sync.RWMutex
rw.RLock() // 读锁(可并发)
// 读操作
rw.RUnlock()
rw.Lock() // 写锁(排他)
// 写操作
rw.Unlock()
// TryLock(Go 1.18+)
if mu.TryLock() {
defer mu.Unlock()
// 获取成功
} else {
// 获取失败,不阻塞
}
| 对比 | Mutex | RWMutex |
|---|
| 读并发 | 否(读也互斥) | 是(多个读可并发) |
| 写并发 | 否 | 否 |
| 适用 | 简单场景 | 读远多于写 |
| 性能 | 简单场景更高 | 读多场景更高 |
Once
var (
once sync.Once
instance *Database
)
func GetDB() *Database {
once.Do(func() {
instance = connectDB() // 只执行一次,即使并发调用
})
return instance
}
Cond
// 条件变量:等待某个条件成立
var (
mu sync.Mutex
cond = sync.NewCond(&mu)
queue []int
)
// 消费者:等待队列非空
func consume() int {
mu.Lock()
for len(queue) == 0 {
cond.Wait() // 释放锁并等待,被唤醒后重新获取锁
}
item := queue[0]
queue = queue[1:]
mu.Unlock()
return item
}
// 生产者
func produce(item int) {
mu.Lock()
queue = append(queue, item)
cond.Signal() // 唤醒一个等待者
// cond.Broadcast() // 唤醒所有等待者
mu.Unlock()
}
Pool
// 对象池:复用对象,减少 GC 压力
var bufPool = sync.Pool{
New: func() interface{} {
return bytes.NewBuffer(make([]byte, 0, 4096))
},
}
func process(data []byte) string {
buf := bufPool.Get().(*bytes.Buffer)
defer func() {
buf.Reset()
bufPool.Put(buf)
}()
buf.Write(data)
// 处理...
return buf.String()
}
并发陷阱
陷阱1:goroutine 泄漏
// ❌ 泄漏:如果 doWork 超时返回,goroutine 永远阻塞在 ch <- result
func leaky() {
ch := make(chan int)
go func() {
result := doWork()
ch <- result // 如果没人接收,永远阻塞
}()
select {
case v := <-ch:
fmt.Println(v)
case <-time.After(time.Second):
return // goroutine 泄漏!
}
}
// ✓ 修复:用缓冲 channel
func fixed() {
ch := make(chan int, 1) // 缓冲为1,发送不阻塞
go func() {
ch <- doWork()
}()
select {
case v := <-ch:
fmt.Println(v)
case <-time.After(time.Second):
return // goroutine 会完成发送后退出
}
}
陷阱2:循环变量捕获
// ❌ Go < 1.22:所有 goroutine 可能共享同一个 i
for i := 0; i < 5; i++ {
go func() {
fmt.Println(i) // 可能全部输出 5
}()
}
// ✓ Go < 1.22 修复方式
for i := 0; i < 5; i++ {
i := i // 创建新变量
go func() {
fmt.Println(i)
}()
}
// ✓ Go >= 1.22:循环变量每次迭代都是新变量,无需修复
// 默认行为已改变,go.mod 中 go 1.22+ 即可
陷阱3:map 并发写
// ❌ panic: concurrent map writes
m := make(map[string]int)
go func() { m["a"] = 1 }()
go func() { m["b"] = 2 }()
// ✓ 用 sync.Map 或 mutex 保护
陷阱4:误用 defer 在循环中
// ❌ defer 在函数返回时才执行,循环中会堆积
for _, file := range files {
f, _ := os.Open(file)
defer f.Close() // 所有 fd 直到函数返回才关闭!
}
// ✓ 在循环内不用 defer,或抽取到子函数
for _, file := range files {
processFile(file) // 子函数内 defer f.Close() 正常
}
func processFile(path string) {
f, err := os.Open(path)
if err != nil {
return
}
defer f.Close()
// 处理
}
并发设计决策
| 需求 | 推荐方案 | 不推荐 |
|---|
| goroutine 间传递数据 | channel | 共享变量 + 锁 |
| 保护共享状态 | Mutex / RWMutex | channel(杀鸡用牛刀) |
| 等待 N 个 goroutine | WaitGroup | channel 手动计数 |
| 一次性初始化 | sync.Once | 双重检查锁 |
| 读多写少的 map | sync.Map | map + RWMutex |
| 对象复用 | sync.Pool | 全局变量 |
| 超时/取消传播 | context | 手动 channel |
| 限流 | 带缓冲 channel | 第三方库(简单场景) |
| 定时任务 | time.Ticker | for + time.Sleep |