并发编程
约 1518 字大约 5 分钟
2026-02-05
Go 语言采用 CSP(Communicating Sequential Processes) 并发模型,提倡通过 **「通信共享内存」**而非 「共享内存实现通信」。
goroutine
Goroutine 是 Go 程序中最基本的并发执行单元。使用 go 关键字即可创建一个 goroutine。
func Hello() {
fmt.Println("Hello Goroutine")
}
func main() {
go Hello() // 开启一个 goroutine
fmt.Println("Hello Main Goroutine")
}注意
main 函数本身是一个 main goroutine,两个 goroutine 并发执行。 若 main goroutine 先结束,Hello goroutine 可能来不及执行完就被终止。
三种可能的输出结果:
Hello Goroutine
Hello Main Goroutine或
Hello Main Goroutine或
Hello Main Goroutine
Hello Goroutinesync.WaitGroup
使用 sync.WaitGroup 等待所有 goroutine 执行完毕后再退出程序。
单个 goroutine
import (
"fmt"
"sync"
)
var wait sync.WaitGroup
func Hello() {
fmt.Println("Hello Goroutine")
wait.Done() // 计数器 -1
}
func main() {
wait.Add(1) // 计数器 +1
go Hello()
fmt.Println("Hello Main Goroutine")
wait.Wait() // 阻塞,等待所有 goroutine 完成
}多个 goroutine
var wait sync.WaitGroup
func Hello(i int) {
fmt.Println("Hello Goroutine", i)
wait.Done()
}
func main() {
wait.Add(10000)
for i := 0; i < 10000; i++ {
go Hello(i)
}
wait.Wait()
}匿名函数 + WaitGroup
for i := 0; i < 10000; i++ {
go func() {
// 注意:这里 i 引用外部 for 循环,形成闭包
fmt.Println("hello", i)
wait.Done()
}()
}Channel(通道)
Channel 是 goroutine 之间的通信机制,遵循**先入先出(FIFO)**规则。
提示
channel 和 slice、map 一样属于引用类型,必须使用 make 初始化。
无缓冲通道(同步通道)
// 无缓冲通道:发送和接收必须同时准备好,否则阻塞
ch := make(chan int)
ch <- 10 // 阻塞,直到有人接收
x := <-ch // 接收有缓冲通道(异步通道)
// 有缓冲通道:容量为1,可以先存再取
ch := make(chan int, 1)
ch <- 10 // ✅ 不阻塞(缓冲区有空位)
fmt.Println("发送成功")
x := <-ch // 取值// 更完整的示例
func main() {
ch := make(chan int, 1) // 创建容量为1的缓冲通道
ch <- 10
x := <-ch
fmt.Println(x) // 10
close(ch)
}通道操作
| 操作 | 语法 | 说明 |
|---|---|---|
| 创建 | make(chan 类型, 容量) | 无缓冲则不指定容量 |
| 发送 | ch <- 值 | 向通道发送数据 |
| 接收 | <- ch 或 x := <-ch | 从通道接收数据 |
| 关闭 | close(ch) | 关闭通道,不再发送 |
| 长度 | len(ch) | 通道中元素数量 |
| 容量 | cap(ch) | 通道缓冲区大小 |
关闭通道后的特点
- 对关闭的通道发送值 →
panic - 对关闭的通道接收值 → 持续获取直到通道为空
- 对关闭且为空的通道接收 → 返回零值 | 4. 关闭已关闭的通道 →
panic
select 多路复用
select 语句用于同时监听多个 channel 的操作,类似于 switch,但专门用于 channel 通信。
select {
case v := <-ch1:
fmt.Println("从 ch1 收到:", v)
case ch2 <- 42:
fmt.Println("向 ch2 发送成功")
case <-ch3:
fmt.Println("从 ch3 收到值")
default:
fmt.Println("所有 channel 都不可用")
}select会阻塞,直到某个 case 可以执行- 如果多个 case 同时就绪,Go 会随机选择一个执行
default分支:当所有 case 都无法执行时,立即执行 default,实现非阻塞操作
超时处理
结合 time.After 可以实现超时控制:
func main() {
ch := make(chan int)
go func() {
time.Sleep(2 * time.Second)
ch <- 1
}()
select {
case v := <-ch:
fmt.Println("收到:", v)
case <-time.After(1 * time.Second):
fmt.Println("超时:未在 1 秒内收到数据")
}
}非阻塞收发
利用 default 分支实现非阻塞的 channel 操作:
func main() {
ch := make(chan int)
select {
case v := <-ch:
fmt.Println("收到:", v)
default:
fmt.Println("ch 没有数据,不阻塞")
}
}for-range 遍历通道
使用 for range 可以持续从 channel 中读取数据,直到 channel 被关闭:
func main() {
ch := make(chan int)
go func() {
for i := 0; i < 5; i++ {
ch <- i
}
close(ch) // 必须关闭,否则 range 会死锁
}()
for v := range ch {
fmt.Println(v) // 0 1 2 3 4
}
}注意
for range会一直阻塞读取,直到 channel 被关闭- 如果 channel 未关闭且没有发送者继续写,会引发死锁(fatal error: all goroutines are asleep - deadlock!)
- 发送方必须在发送完毕后调用
close(ch)来通知接收方结束
互斥锁 sync.Mutex
当多个 goroutine 同时读写共享内存时,需要使用互斥锁来保护临界区,防止数据竞争。
sync.Mutex(互斥锁)
var (
x int
mu sync.Mutex
wait sync.WaitGroup
)
func AddWithLock() {
for i := 0; i < 5000; i++ {
mu.Lock() // 加锁
x++ // 临界区:同时只能有一个 goroutine 执行
mu.Unlock() // 解锁
}
wait.Done()
}
func main() {
wait.Add(2)
go AddWithLock()
go AddWithLock()
wait.Wait()
fmt.Println(x) // 10000 —— 结果正确
}注意
不加锁时,x++ 并非原子操作(实际是读取→修改→写入三步),多个 goroutine 并发执行会导致数据竞争(data race),最终结果会小于预期值。
sync.RWMutex(读写锁)
sync.RWMutex 是对 Mutex 的优化,适用于读多写少的场景:
| 操作 | 方法 | 并发规则 |
|---|---|---|
| 读锁 | RLock() / RUnlock() | 多个读可并发,写必须等待所有读完成 |
| 写锁 | Lock() / Unlock() | 独占,阻塞所有读写 |
var (
data map[string]string
rwMu sync.RWMutex
)
func Read(key string) string {
rwMu.RLock() // 加读锁
defer rwMu.RUnlock() // 解锁
return data[key]
}
func Write(key, value string) {
rwMu.Lock() // 加写锁
defer rwMu.Unlock() // 解锁
data[key] = value
}Worker Pool 模式
Worker Pool(工作池)是一种经典的并发模式:启动固定数量的 worker goroutine,从 job channel 领取任务,将结果写入 results channel。
func worker(id int, jobs <-chan int, results chan<- int, wg *sync.WaitGroup) {
defer wg.Done()
for j := range jobs {
fmt.Printf("worker %d 开始任务 %d\n", id, j)
time.Sleep(time.Second) // 模拟耗时操作
results <- j * 2
fmt.Printf("worker %d 完成任务 %d\n", id, j)
}
}
func main() {
const numJobs = 10
const numWorkers = 3
jobs := make(chan int, numJobs)
results := make(chan int, numJobs)
// 启动 worker
var wg sync.WaitGroup
for w := 1; w <= numWorkers; w++ {
wg.Add(1)
go worker(w, jobs, results, &wg)
}
// 发送任务
for j := 1; j <= numJobs; j++ {
jobs <- j
}
close(jobs) // 关闭 job 通道,通知 worker 不再有新任务
// 等待所有 worker 完成
wg.Wait()
close(results) // 关闭 results 通道
// 收集结果
for r := range results {
fmt.Println("结果:", r)
}
}模式要点
| 要点 | 说明 |
|---|---|
| 固定 Worker | 控制并发数量,防止 goroutine 无限膨胀 |
| Job Channel | 使用 chan<-(只写方向)限制 worker 仅接收 |
| Result Channel | 使用 <-chan(只读方向)隔离职责 |
| 关闭通道 | 发送完毕后 close(jobs),worker 的 for-range 自然退出 |
| WaitGroup | 主 goroutine 等待所有 worker 结束后再关闭 results |
