Skip to content

并发编程

本章讲解Go并发编程:Goroutine轻量级线程、Channel通信(CSP模型)、同步原语(WaitGroup/Mutex)、select多路复用、竞态检测。重点:理解”通过通信共享内存”核心理念,掌握channel和mutex的适用场景。

前置知识:Go基础语法(函数、struct) 学习目标:能编写并发程序、理解GMP调度模型、避免数据竞争


本节目的:对比Python并发模型,理解Goroutine的轻量与高效

# 线程(受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限制真正的并行
  • 线程开销大
  • 异步代码复杂

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
创建成本~1MB栈~2KB栈(可增长)
创建速度慢极快
切换成本上下文切换极低
调度OS内核Go运行时(GMP)
数量限制受内存限制可轻松创建百万
// 匿名函数
go func() {
fmt.Println("在goroutine中运行")
}()
// 传递参数
go func(msg string) {
fmt.Println(msg)
}("hello")

主函数返回,所有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("主函数退出")
}

本节目的:理解CSP模型、掌握channel用法(创建/发送/接收/关闭)

Go并发哲学:“不要通过共享内存来通信;相反,通过通信来共享内存”

Channel是goroutine之间的通信通道。

CSP图示:

┌─────────────┐ chan ┌─────────────┐
│ Goroutine │ ────────────→ │ Goroutine │
│ (Producer) │ send(v) │ (Consumer) │
└─────────────┘ ←──────────── └─────────────┘
close(ch)
  • Producer通过channel发送数据
  • Consumer通过channel接收数据
  • 无需共享内存,天然避免数据竞争
// 创建无缓冲channel
ch := make(chan int)
// 创建有缓冲channel
ch := make(chan int, 10)
ch := make(chan int)
// 发送
go func() {
ch <- 42 // 发送数据
}()
// 接收
value := <-ch // 从channel接收
fmt.Println(value) // 42
// 关闭channel
close(ch)

发送和接收必须同时准备:

ch := make(chan int)
// 发送goroutine
go func() {
fmt.Println("发送前")
ch <- 100
fmt.Println("发送后")
}()
// 接收goroutine
fmt.Println("接收前:", <-ch)
fmt.Println("接收后")

不阻塞直到缓冲区满:

ch := make(chan int, 3)
// 发送3次不阻塞
ch <- 1
ch <- 2
ch <- 3
// 再发送会阻塞(缓冲区满)
// ch <- 4 // 阻塞
// 接收
fmt.Println(<-ch) // 1
fmt.Println(<-ch) // 2
ch := make(chan string, 2)
// 发送
ch <- "hello"
ch <- "world"
// 接收
msg1 := <-ch
msg2 := <-ch
// 检查channel状态
value, ok := <-ch // ok为false表示channel已关闭且无数据
// 关闭
close(ch)
// range遍历channel
for msg := range ch {
fmt.Println(msg)
}
// 生产者:只能发送
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)

本节目的:掌握WaitGroup等待多goroutine、Mutex保护共享资源

等待一组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")

保护共享资源:

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
}

读多写少场景:

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
}

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()
}
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
// ...

本节目的:掌握select监听多channel、实现超时和优雅退出

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 {
case v := <-ch:
fmt.Println("收到:", v)
default:
fmt.Println("没有数据")
}
// 如果channel不可用,立即执行default
ch := 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)
}

本节目的:理解数据竞争、掌握-race检测工具、学会三种解决方案

两个goroutine同时访问同一内存,至少一个写入:

var counter int
func increment() {
counter++ // 竞态!
}
func main() {
for i := 0; i < 1000; i++ {
go increment()
}
time.Sleep(time.Second)
fmt.Println(counter) // 不确定的值
}

Go内置竞态检测器:

Terminal window
go run -race main.go
go test -race ./...
go build -race

会检测到:

WARNING: DATA RACE
Read at 0x00c000100008 by goroutine 7:
...
Previous write at 0x00c000100008 by goroutine 5:
...

方法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)
}
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 int32
atomic.StoreInt32(&flag, 1)
atomic.CompareAndSwapInt32(&flag, 1, 0)
// 指针原子操作
atomic.LoadPointer(&ptrVal)
atomic.StorePointer(&ptrVal, new(int))

  1. 优先选择channel:适合逻辑上的数据流传递(生产者-消费者、流水线)
  2. 互斥锁保护复杂状态:当状态涉及多个字段、需要复杂条件判断时用Mutex
  3. 避免嵌套锁:轻易导致死锁;优先设计无锁结构或单层锁
  4. 关闭channel的原则:发送方关闭(close),接收方通过ok判断channel是否已关闭
  5. select+context:结合context实现优雅取消,避免 goroutine 泄漏
  6. 竞态检测是必须的:上线前用go run -race测试,检测隐藏的数据竞争

PythonGo
threading.Threadgo func(){}
threading.Locksync.Mutex
threading.Semaphoresync.WaitGroup
queue.Queuechan
asyncio.gathergoroutine + channel
无select

Go并发核心概念:

  • goroutine:轻量级并发
  • channel:通信媒介
  • select:多路复用
  • sync包:同步原语

  1. 启动10个goroutine,用WaitGroup等待它们完成
  2. 创建有缓冲channel,实现生产者-消费者模式
  3. 用select实现超时机制
  4. 用互斥锁保护map的读写
  5. 用-race选项运行竞态检测,修复数据竞争