主题
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 会被强制终止,因此需通过channel或sync.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()— 根 contextcontext.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 会持续占用资源,注意通过
context或donechannel 回收。
作者:yanshaodong