并发编程
本章讲解Go并发编程:Goroutine轻量级线程、Channel通信(CSP模型)、同步原语(WaitGroup/Mutex)、select多路复用、竞态检测。重点:理解”通过通信共享内存”核心理念,掌握channel和mutex的适用场景。
前置知识:Go基础语法(函数、struct) 学习目标:能编写并发程序、理解GMP调度模型、避免数据竞争
10.1 从多线程到Goroutine
Section titled “10.1 从多线程到Goroutine”本节目的:对比Python并发模型,理解Goroutine的轻量与高效
Python的并发模型
Section titled “Python的并发模型”# 线程(受GIL限制)import threading
def worker(n): print(f"Worker {n}")
threads = []for i in range(5): t = threading.Thread(target=worker, args=(i,)) threads.append(t) t.start()
for t in threads: t.join()# 异步(协程)import asyncio
async def worker(n): print(f"Worker {n}")
async def main(): await asyncio.gather(*[worker(i) for i in range(5)])
asyncio.run(main())Python的问题:
- GIL限制真正的并行
- 线程开销大
- 异步代码复杂
Go的Goroutine
Section titled “Go的Goroutine”goroutine是轻量级线程,由Go运行时管理:
package main
import ( "fmt" "time")
func worker(n int) { fmt.Printf("Worker %d\n", n)}
func main() { // 启动5个goroutine for i := 0; i < 5; i++ { go worker(i) }
// 等待goroutine完成 time.Sleep(time.Second) fmt.Println("Done")}Goroutine vs 线程
Section titled “Goroutine vs 线程”| 特性 | 线程 | Goroutine |
|---|---|---|
| 创建成本 | ~1MB栈 | ~2KB栈(可增长) |
| 创建速度 | 慢 | 极快 |
| 切换成本 | 上下文切换 | 极低 |
| 调度 | OS内核 | Go运行时(GMP) |
| 数量限制 | 受内存限制 | 可轻松创建百万 |
启动多个goroutine
Section titled “启动多个goroutine”// 匿名函数go func() { fmt.Println("在goroutine中运行")}()
// 传递参数go func(msg string) { fmt.Println(msg)}("hello")主goroutine退出
Section titled “主goroutine退出”主函数返回,所有goroutine终止:
func main() { go func() { time.Sleep(time.Second) fmt.Println("这个不会打印") }()
fmt.Println("主函数退出") // 立即退出}
// 输出: 主函数退出解决方案:等待
import "sync"
func main() { var wg sync.WaitGroup wg.Add(1)
go func() { defer wg.Done() fmt.Println("goroutine完成") }()
wg.Wait() // 等待 fmt.Println("主函数退出")}10.2 Channel:通信而非共享内存
Section titled “10.2 Channel:通信而非共享内存”本节目的:理解CSP模型、掌握channel用法(创建/发送/接收/关闭)
Go并发哲学:“不要通过共享内存来通信;相反,通过通信来共享内存”
Channel是goroutine之间的通信通道。
CSP图示:
┌─────────────┐ chan ┌─────────────┐│ Goroutine │ ────────────→ │ Goroutine ││ (Producer) │ send(v) │ (Consumer) │└─────────────┘ ←──────────── └─────────────┘ close(ch)- Producer通过channel发送数据
- Consumer通过channel接收数据
- 无需共享内存,天然避免数据竞争
创建Channel
Section titled “创建Channel”// 创建无缓冲channelch := make(chan int)
// 创建有缓冲channelch := make(chan int, 10)ch := make(chan int)
// 发送go func() { ch <- 42 // 发送数据}()
// 接收value := <-ch // 从channel接收fmt.Println(value) // 42
// 关闭channelclose(ch)无缓冲Channel
Section titled “无缓冲Channel”发送和接收必须同时准备:
ch := make(chan int)
// 发送goroutinego func() { fmt.Println("发送前") ch <- 100 fmt.Println("发送后")}()
// 接收goroutinefmt.Println("接收前:", <-ch)fmt.Println("接收后")有缓冲Channel
Section titled “有缓冲Channel”不阻塞直到缓冲区满:
ch := make(chan int, 3)
// 发送3次不阻塞ch <- 1ch <- 2ch <- 3
// 再发送会阻塞(缓冲区满)// ch <- 4 // 阻塞
// 接收fmt.Println(<-ch) // 1fmt.Println(<-ch) // 2Channel基本操作
Section titled “Channel基本操作”ch := make(chan string, 2)
// 发送ch <- "hello"ch <- "world"
// 接收msg1 := <-chmsg2 := <-ch
// 检查channel状态value, ok := <-ch // ok为false表示channel已关闭且无数据
// 关闭close(ch)
// range遍历channelfor msg := range ch { fmt.Println(msg)}单向Channel
Section titled “单向Channel”// 生产者:只能发送func producer(ch chan<- int) { for i := 0; i < 5; i++ { ch <- i } close(ch)}
// 消费者:只能接收func consumer(ch <-chan int) { for v := range ch { fmt.Println(v) }}
// 使用ch := make(chan int)go producer(ch)consumer(ch)10.3 同步原语:WaitGroup / Mutex
Section titled “10.3 同步原语:WaitGroup / Mutex”本节目的:掌握WaitGroup等待多goroutine、Mutex保护共享资源
WaitGroup
Section titled “WaitGroup”等待一组goroutine完成:
import "sync"
func main() { var wg sync.WaitGroup
// 添加计数 wg.Add(3)
for i := 0; i < 3; i++ { go func(id int) { defer wg.Done() // 完成时递减 fmt.Printf("Worker %d done\n", id) }(i) }
// 等待所有goroutine wg.Wait() fmt.Println("All workers done")}对比Python:
import threading
wg = threading.Semaphore(0) # 初始值0
def worker(id): print(f"Worker {id} done") wg.release() # 递增
for i in range(3): threading.Thread(target=worker, args=(i,)).start()
wg.acquire(3) # 等待3次print("All workers done")Mutex(互斥锁)
Section titled “Mutex(互斥锁)”保护共享资源:
import "sync"
var ( counter int mu sync.Mutex)
func increment() { mu.Lock() // 加锁 defer mu.Unlock() // 释放 counter++}
func main() { var wg sync.WaitGroup wg.Add(100)
for i := 0; i < 100; i++ { go func() { increment() wg.Done() }() }
wg.Wait() fmt.Println(counter) // 100}RWMutex(读写锁)
Section titled “RWMutex(读写锁)”读多写少场景:
var ( mu sync.RWMutex data map[string]string)
func read(key string) string { mu.RLock() defer mu.RUnlock() return data[key]}
func write(key, value string) { mu.Lock() defer mu.Unlock() data[key] = value}Cond(条件变量)
Section titled “Cond(条件变量)”goroutine等待特定条件:
var ( mu sync.Mutex cond = sync.NewCond(&mu) ready bool)
func worker() { mu.Lock() for !ready { cond.Wait() // 等待条件 } mu.Unlock() fmt.Println("Worker 开始工作")}
func main() { go worker()
// 准备就绪 mu.Lock() ready = true cond.Signal() // 通知一个goroutine mu.Unlock()}Once(只执行一次)
Section titled “Once(只执行一次)”var ( once sync.Once data string)
func getData() string { once.Do(func() { data = "initialized" fmt.Println("初始化一次") }) return data}
func main() { for i := 0; i < 5; i++ { go func() { fmt.Println(getData()) }() } time.Sleep(time.Second)}
// 输出: 初始化一次// initialized// initialized// ...10.4 select多路复用
Section titled “10.4 select多路复用”本节目的:掌握select监听多channel、实现超时和优雅退出
select基础
Section titled “select基础”select监听多个channel操作:
ch1 := make(chan int)ch2 := make(chan int)
go func() { time.Sleep(100 * time.Millisecond) ch1 <- 1}()
go func() { time.Sleep(50 * time.Millisecond) ch2 <- 2}()
// 哪个先完成就处理哪个select {case v := <-ch1: fmt.Println("ch1:", v)case v := <-ch2: fmt.Println("ch2:", v)}// 输出: ch2: 2(ch2先完成)select与default
Section titled “select与default”select {case v := <-ch: fmt.Println("收到:", v)default: fmt.Println("没有数据")}// 如果channel不可用,立即执行defaultch := make(chan int)
select {case v := <-ch: fmt.Println("收到:", v)case <-time.After(time.Second): fmt.Println("超时")}for { select { case v := <-ch: fmt.Println("收到:", v) // 无default,会永远阻塞直到有数据 }}优雅退出goroutine:
func worker(stop <-chan struct{}) { for { select { case <-stop: fmt.Println("收到停止信号") return default: // 正常工作 time.Sleep(100 * time.Millisecond) } }}
func main() { stop := make(chan struct{})
go worker(stop)
time.Sleep(500 * time.Millisecond) close(stop) // 发送停止信号 time.Sleep(100 * time.Millisecond)}10.5 竞态检测与数据竞争
Section titled “10.5 竞态检测与数据竞争”本节目的:理解数据竞争、掌握
-race检测工具、学会三种解决方案
什么是数据竞争
Section titled “什么是数据竞争”两个goroutine同时访问同一内存,至少一个写入:
var counter int
func increment() { counter++ // 竞态!}
func main() { for i := 0; i < 1000; i++ { go increment() } time.Sleep(time.Second) fmt.Println(counter) // 不确定的值}-race检测
Section titled “-race检测”Go内置竞态检测器:
go run -race main.gogo test -race ./...go build -race会检测到:
WARNING: DATA RACERead at 0x00c000100008 by goroutine 7: ... Previous write at 0x00c000100008 by goroutine 5: ...解决竞态的方法
Section titled “解决竞态的方法”方法1:Mutex
var ( counter int mu sync.Mutex)
func increment() { mu.Lock() defer mu.Unlock() counter++}方法2:Atomic
import "sync/atomic"
var counter int64
func increment() { atomic.AddInt64(&counter, 1)}
func main() { var wg sync.WaitGroup for i := 0; i < 1000; i++ { wg.Add(1) go func() { increment() wg.Done() }() } wg.Wait() fmt.Println(atomic.LoadInt64(&counter))}方法3:Channel
func main() { ch := make(chan int, 1) // 缓冲channel作为互斥 var counter int
for i := 0; i < 1000; i++ { go func() { ch <- 1 // 获取"锁" counter++ <-ch // 释放"锁" }() } time.Sleep(time.Second) fmt.Println(counter)}atomic包支持的原子操作
Section titled “atomic包支持的原子操作”import "sync/atomic"
var ( int64Val int64 ptrVal *int)
// 整数原子操作atomic.AddInt64(&int64Val, 1)atomic.LoadInt64(&int64Val)atomic.StoreInt64(&int64Val, 100)atomic.SwapInt64(&int64Val, 200)
// 布尔原子操作var flag int32atomic.StoreInt32(&flag, 1)atomic.CompareAndSwapInt32(&flag, 1, 0)
// 指针原子操作atomic.LoadPointer(&ptrVal)atomic.StorePointer(&ptrVal, new(int))- 优先选择channel:适合逻辑上的数据流传递(生产者-消费者、流水线)
- 互斥锁保护复杂状态:当状态涉及多个字段、需要复杂条件判断时用Mutex
- 避免嵌套锁:轻易导致死锁;优先设计无锁结构或单层锁
- 关闭channel的原则:发送方关闭(close),接收方通过
ok判断channel是否已关闭 - select+context:结合context实现优雅取消,避免 goroutine 泄漏
- 竞态检测是必须的:上线前用
go run -race测试,检测隐藏的数据竞争
| Python | Go |
|---|---|
threading.Thread | go func(){} |
threading.Lock | sync.Mutex |
threading.Semaphore | sync.WaitGroup |
queue.Queue | chan |
asyncio.gather | goroutine + channel |
| 无 | select |
Go并发核心概念:
- goroutine:轻量级并发
- channel:通信媒介
- select:多路复用
- sync包:同步原语
- 启动10个goroutine,用WaitGroup等待它们完成
- 创建有缓冲channel,实现生产者-消费者模式
- 用select实现超时机制
- 用互斥锁保护map的读写
- 用-race选项运行竞态检测,修复数据竞争