Skip to content

Go 并发编程

Go 的并发模型基于 CSP(Communicating Sequential Processes) 思想,核心是两个原语:

  • goroutine:轻量级线程,由 Go runtime 调度,初始栈仅 2KB,可轻松创建百万级。
  • channel:goroutine 之间通信的管道,遵循「不要通过共享内存来通信,而要通过通信来共享内存」。

goroutine

go
func say(s string) {
    for i := 0; i < 3; i++ {
        time.Sleep(100 * time.Millisecond)
        fmt.Println(s)
    }
}

func main() {
    go say("world")   // 启动一个 goroutine,立即返回
    say("hello")      // 主 goroutine 继续执行
}

go 关键字后跟函数调用即可启动 goroutine。主 goroutine 退出时,其他 goroutine 会被强制终止,因此需通过 channelsync.WaitGroup 等待它们完成。

channel

基础用法

go
// 无缓冲 channel:发送和接收必须同时就绪(同步)
ch := make(chan int)

go func() {
    ch <- 42        // 发送,会阻塞直到有人接收
}()

val := <-ch         // 接收,会阻塞直到有人发送
fmt.Println(val)    // 42

// 关闭 channel
close(ch)
// 接收时判断是否已关闭
v, ok := <-ch

带缓冲 channel

go
ch := make(chan int, 100)   // 容量 100,满之前发送不阻塞
ch <- 1
ch <- 2

单向 channel

go
func producer(out chan<- int) {   // 只写
    out <- 1
}
func consumer(in <-chan int) {    // 只读
    fmt.Println(<-in)
}

select 多路复用

select 等待多个 channel 操作,类似 IO 多路复用:

go
func fibonacci(c, quit chan int) {
    x, y := 0, 1
    for {
        select {
        case c <- x:
            x, y = y, x+y
        case <-quit:
            fmt.Println("quit")
            return
        // 默认分支(非阻塞)
        default:
            fmt.Println("no activity")
        // 超时控制
        case <-time.After(2 * time.Second):
            fmt.Println("timeout")
            return
        }
    }
}

sync 同步原语

WaitGroup

等待一组 goroutine 完成:

go
var wg sync.WaitGroup

for i := 0; i < 10; i++ {
    wg.Add(1)
    go func(id int) {
        defer wg.Done()   // 完成时减一
        // 业务逻辑
    }(i)
}
wg.Wait()   // 阻塞直到计数归零

Mutex / RWMutex

保护共享资源:

go
var mu sync.Mutex
var counter int

for i := 0; i < 1000; i++ {
    go func() {
        mu.Lock()
        defer mu.Unlock()
        counter++
    }()
}

// 读写锁:读多写少场景性能更好
var rw sync.RWMutex
rw.RLock()      // 多个读可并发
rw.RUnlock()
rw.Lock()       // 写独占
rw.Unlock()

Once / Map / Pool

go
// sync.Once:保证只执行一次(常用于单例初始化)
var once sync.Once
once.Do(initConfig)

// sync.Map:并发安全的 map(适合读多写少)
var m sync.Map
m.Store("key", "value")
v, ok := m.Load("key")

// sync.Pool:对象复用,降低 GC 压力
var bufPool = sync.Pool{
    New: func() any { return new(bytes.Buffer) },
}
buf := bufPool.Get().(*bytes.Buffer)
buf.Reset()
bufPool.Put(buf)

context 取消传播

context 用于在 goroutine 树中传递取消信号、超时和请求范围的值:

go
func worker(ctx context.Context) {
    for {
        select {
        case <-ctx.Done():            // 收到取消信号
            fmt.Println("canceled:", ctx.Err())
            return
        default:
            // 正常工作
        }
    }
}

func main() {
    // 3 秒后自动取消
    ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
    defer cancel()

    go worker(ctx)
    time.Sleep(5 * time.Second)
}

常用派生函数:

  • context.Background() — 根 context
  • context.WithCancel(parent) — 手动取消
  • context.WithTimeout(parent, d) — 超时取消
  • context.WithDeadline(parent, t) — 截止时间取消
  • context.WithValue(parent, key, val) — 传递请求范围的值(谨慎使用)

约定:context 作为函数第一个参数,命名为 ctx

常见并发模式

Worker Pool(工作池)

go
func worker(id int, jobs <-chan int, results chan<- int) {
    for j := range jobs {          // channel 关闭时循环退出
        results <- j * 2
    }
}

jobs := make(chan int, 100)
results := make(chan int, 100)

for w := 1; w <= 3; w++ {
    go worker(w, jobs, results)
}

for j := 1; j <= 9; j++ {
    jobs <- j
}
close(jobs)

for a := 1; a <= 9; a++ {
    <-results
}

扇出 / 扇入(Fan-out / Fan-in)

多个 goroutine 处理同一 channel(扇出),再将结果汇聚到一个 channel(扇入)。

注意事项

  • 不要通过共享内存通信,优先用 channel。
  • channel 发送/接收在 goroutine 泄漏时会导致永久阻塞,务必保证有接收方。
  • 使用 -race 检测数据竞争:go test -race ./...
  • goroutine 泄漏:未正确退出的 goroutine 会持续占用资源,注意通过 contextdone channel 回收。

作者:yanshaodong