文章

并发模型

并发模型

Go 的并发哲学:不要通过共享内存来通信,而要通过通信来共享内存。

核心组件

组件作用类比
goroutine轻量级线程绿色线程 / 协程
channelgoroutine 间通信管道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 MB2 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)panicpanic成功
长度 len(ch)00当前缓冲元素数
容量 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 {
    // 获取失败,不阻塞
}
对比MutexRWMutex
读并发否(读也互斥)是(多个读可并发)
写并发
适用简单场景读远多于写
性能简单场景更高读多场景更高

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 / RWMutexchannel(杀鸡用牛刀)
等待 N 个 goroutineWaitGroupchannel 手动计数
一次性初始化sync.Once双重检查锁
读多写少的 mapsync.Mapmap + RWMutex
对象复用sync.Pool全局变量
超时/取消传播context手动 channel
限流带缓冲 channel第三方库(简单场景)
定时任务time.Tickerfor + time.Sleep