Go 并发编程

4 min高级

goroutine 原理、channel、select、sync 包、context 包、并发模式与竞态检测。

前置知识

学习目标

  • 掌握「1. Goroutine」的核心机制、典型用法与常见陷阱
  • 掌握「2. Channel」的核心机制、典型用法与常见陷阱
  • 掌握「3. Select」的核心机制、典型用法与常见陷阱
  • 掌握「4. sync 包」的核心机制、典型用法与常见陷阱
  • 掌握「5. Context 包」的核心机制、典型用法与常见陷阱

1. Goroutine

1.1 基本使用

Goroutine 是 Go 运行时管理的轻量级线程,由 go 关键字启动:

// 启动 goroutine
go func() {
    fmt.Println("并发执行")
}()

// 启动函数
go doWork()

// 主 goroutine 不会等待子 goroutine
func main() {
    go fmt.Println("可能看不到这行")
    // main 退出,所有 goroutine 终止
}

1.2 Goroutine vs OS 线程

特性GoroutineOS 线程
初始栈大小2KB(可动态伸缩)1-8MB(固定)
创建成本微秒级毫秒级
调度Go 运行时(M:N 模型)操作系统内核(1:1)
切换成本~100ns(用户态)~1-10μs(内核态)
数量上限百万级千级
通信channel共享内存 + 锁

1.3 GMP 调度模型

flowchart TD
    G[G goroutine 协程,用户级轻量线程]
    M[M machine 操作系统线程]
    PP[P processor 逻辑处理器,持有本地运行队列]
    S[Scheduler]
    S --> P0[P0 [G G]]
    S --> P1[P1 [G G]]
    S --> P2[P2 [G G]]
    S --> P3[P3 [G G]]
    S --> GQ[全局队列 [G G G]]
    P0 --> M0[M0]
    P1 --> M1[M1]
    P2 --> M2[M2]
    P3 --> M3[M3]

调度策略:

  • Work Stealing:P 的本地队列为空时,从其他 P 或全局队列窃取 G
  • Hand Off:M 阻塞(如系统调用)时,P 绑定到新的 M 继续运行
  • 抢占式调度:基于协作(函数调用检查)+ 基于信号(Go 1.14+)

1.4 等待 Goroutine 完成

// 使用 WaitGroup
func main() {
    var wg sync.WaitGroup

    for i := 0; i < 5; i++ {
        wg.Add(1)
        go func(id int) {
            defer wg.Done()
            fmt.Printf("Worker %d done\n", id)
        }(i)
    }

    wg.Wait()
    fmt.Println("All workers finished")
}

2. Channel

2.1 基本操作

// 无缓冲 channel(同步通道)
ch := make(chan int)

// 有缓冲 channel
ch := make(chan int, 100)

// 发送
ch <- 42

// 接收
v := <-ch

// 接收并检查是否关闭
v, ok := <-ch
if !ok {
    fmt.Println("channel 已关闭")
}

// 关闭 channel
close(ch)

// 遍历 channel(直到关闭)
for v := range ch {
    fmt.Println(v)
}

2.2 无缓冲 vs 有缓冲

// 无缓冲:发送和接收必须同时就绪(同步)
ch := make(chan int)
go func() {
    ch <- 1 // 阻塞直到有人接收
}()
v := <-ch  // 阻塞直到有数据

// 有缓冲:缓冲区满前发送不阻塞
ch := make(chan int, 3)
ch <- 1 // 不阻塞
ch <- 2 // 不阻塞
ch <- 3 // 不阻塞
// ch <- 4 // 阻塞(缓冲区满)

2.3 单向 Channel

// 只发送
func producer(ch chan<- int) {
    for i := 0; i < 10; i++ {
        ch <- i
    }
    close(ch)
}

// 只接收
func consumer(ch <-chan int) {
    for v := range ch {
        fmt.Println(v)
    }
}

func main() {
    ch := make(chan int, 5)
    go producer(ch)
    consumer(ch)
}

2.4 Channel 底层结构

flowchart TD
    H[hchan 结构]
    H --> F1[buf *array 环形缓冲区]
    H --> F2[sendx uint 发送索引]
    H --> F3[recvx uint 接收索引]
    H --> F4[qcount uint 缓冲区元素数]
    H --> F5[dataqsiz uint 缓冲区大小]
    H --> F6[elemtype *type 元素类型]
    H --> F7[closed uint32 是否关闭]
    H --> F8[sendq waitq 发送等待队列]
    H --> F9[recvq waitq 接收等待队列]
    H --> F10[lock mutex 互斥锁]

3. Select

select 同时监听多个 channel 操作:

select {
case v := <-ch1:
    fmt.Println("ch1:", v)
case v := <-ch2:
    fmt.Println("ch2:", v)
case ch3 <- 42:
    fmt.Println("sent to ch3")
default:
    fmt.Println("没有就绪的 channel")
}

3.1 超时控制

select {
case result := <-ch:
    fmt.Println("收到结果:", result)
case <-time.After(5 * time.Second):
    fmt.Println("超时")
}

3.2 非阻塞操作

select {
case msg := <-ch:
    fmt.Println(msg)
default:
    // channel 无数据时立即执行
    fmt.Println("无数据")
}

3.3 退出信号

func worker(done <-chan struct{}) {
    for {
        select {
        case <-done:
            fmt.Println("收到退出信号")
            return
        default:
            doWork()
        }
    }
}

4. sync 包

4.1 Mutex(互斥锁)

type SafeCounter struct {
    mu sync.Mutex
    m  map[string]int
}

func (c *SafeCounter) Inc(key string) {
    c.mu.Lock()
    defer c.mu.Unlock()
    c.m[key]++
}

func (c *SafeCounter) Get(key string) int {
    c.mu.Lock()
    defer c.mu.Unlock()
    return c.m[key]
}

4.2 RWMutex(读写锁)

type Cache struct {
    mu   sync.RWMutex
    data map[string]string
}

func (c *Cache) Get(key string) (string, bool) {
    c.mu.RLock()         // 读锁,允许多个并发读
    defer c.mu.RUnlock()
    v, ok := c.data[key]
    return v, ok
}

func (c *Cache) Set(key, value string) {
    c.mu.Lock()          // 写锁,独占
    defer c.mu.Unlock()
    c.data[key] = value
}

4.3 WaitGroup

func fetchAll(urls []string) []string {
    var (
        wg    sync.WaitGroup
        mu    sync.Mutex
        results []string
    )

    for _, url := range urls {
        wg.Add(1)
        go func(u string) {
            defer wg.Done()
            data := fetch(u)
            mu.Lock()
            results = append(results, data)
            mu.Unlock()
        }(url)
    }

    wg.Wait()
    return results
}

4.4 Once

var (
    instance *Config
    once     sync.Once
)

func GetConfig() *Config {
    once.Do(func() {
        instance = loadConfig() // 只执行一次
    })
    return instance
}

Go 1.21 起标准库补充了三个基于 sync.Once 的泛型便捷函数,一行完成”只算一次 + 缓存结果”:

// sync.OnceFunc:包装一次初始化的函数(Go 1.21+)
var ready = sync.OnceFunc(loadConfig) // loadConfig 只会被调用一次

// sync.OnceValue:带返回值的版本,首次调用结果被缓存
var defaultPrice = sync.OnceValue(func() int {
    return queryPriceFromDB() // 后续调用直接返回缓存值
})

4.5 sync.Map

并发安全的 map,适用于读多写少场景:

var m sync.Map

// 存储
m.Store("key", "value")

// 读取
v, ok := m.Load("key")

// 读取或写入(原子操作)
actual, loaded := m.LoadOrStore("key", "default")

// 删除
m.Delete("key")

// 遍历
m.Range(func(key, value any) bool {
    fmt.Println(key, value)
    return true // 返回 false 停止遍历
})

4.6 sync.Pool

对象池,减少 GC 压力:

var bufPool = sync.Pool{
    New: func() any {
        return new(bytes.Buffer)
    },
}

func process(data []byte) {
    buf := bufPool.Get().(*bytes.Buffer)
    defer func() {
        buf.Reset()
        bufPool.Put(buf)
    }()

    buf.Write(data)
    // 使用 buf...
}

5. Context 包

Context 用于在 goroutine 之间传递取消信号、超时和值:

5.1 创建 Context

// 根 context(不可取消)
ctx := context.Background()
ctx := context.TODO() // 不确定用哪个时使用

// 可取消
ctx, cancel := context.WithCancel(ctx)
defer cancel() // 确保资源释放

// 超时取消
ctx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()

// 截止时间取消
ctx, cancel := context.WithDeadline(ctx, time.Now().Add(10*time.Second))
defer cancel()

// 传递值
ctx = context.WithValue(ctx, "requestID", "abc-123")

5.2 使用 Context

func fetchData(ctx context.Context, url string) ([]byte, error) {
    req, err := http.NewRequestWithContext(ctx, "GET", url, nil)
    if err != nil {
        return nil, err
    }
    resp, err := http.DefaultClient.Do(req)
    if err != nil {
        return nil, err
    }
    defer resp.Body.Close()
    return io.ReadAll(resp.Body)
}

// 超时控制
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
data, err := fetchData(ctx, "https://api.example.com/data")
if err != nil {
    if ctx.Err() == context.DeadlineExceeded {
        fmt.Println("请求超时")
    }
}

5.3 传播取消

func handler(ctx context.Context) {
    ctx, cancel := context.WithTimeout(ctx, 10*time.Second)
    defer cancel()

    // 启动子任务
    results := make(chan string, 3)
    for i := 0; i < 3; i++ {
        go func(id int) {
            result, err := doTask(ctx, id)
            if err != nil {
                cancel() // 任一任务失败,取消所有任务
                return
            }
            results <- result
        }(i)
    }

    for i := 0; i < 3; i++ {
        select {
        case r := <-results:
            fmt.Println(r)
        case <-ctx.Done():
            fmt.Println("取消:", ctx.Err())
            return
        }
    }
}

6. 并发模式

6.1 Fan-in / Fan-out

// Fan-out:将工作分发到多个 goroutine
func fanOut(input <-chan int, n int) []<-chan int {
    channels := make([]<-chan int, n)
    for i := 0; i < n; i++ {
        channels[i] = worker(input)
    }
    return channels
}

func worker(input <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for v := range input {
            out <- process(v)
        }
    }()
    return out
}

// Fan-in:合并多个 channel
func fanIn(channels ...<-chan int) <-chan int {
    out := make(chan int)
    var wg sync.WaitGroup
    for _, ch := range channels {
        wg.Add(1)
        go func(c <-chan int) {
            defer wg.Done()
            for v := range c {
                out <- v
            }
        }(ch)
    }
    go func() {
        wg.Wait()
        close(out)
    }()
    return out
}

6.2 Pipeline

func generate(nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for _, n := range nums {
            out <- n
        }
    }()
    return out
}

func square(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            out <- n * n
        }
    }()
    return out
}

func filter(in <-chan int, pred func(int) bool) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            if pred(n) {
                out <- n
            }
        }
    }()
    return out
}

// 链式调用
pipeline := filter(square(generate(1, 2, 3, 4, 5)), func(n int) bool {
    return n > 10
})
for v := range pipeline {
    fmt.Println(v) // 16, 25
}

6.3 Worker Pool

func workerPool(ctx context.Context, jobs <-chan Job, results chan<- Result, n int) {
    var wg sync.WaitGroup
    for i := 0; i < n; i++ {
        wg.Add(1)
        go func(id int) {
            defer wg.Done()
            for {
                select {
                case job, ok := <-jobs:
                    if !ok {
                        return
                    }
                    results <- processJob(ctx, id, job)
                case <-ctx.Done():
                    return
                }
            }
        }(i)
    }
    go func() {
        wg.Wait()
        close(results)
    }()
}

6.4 限流器

// 令牌桶限流
func rateLimiter(ctx context.Context, interval time.Duration) <-chan struct{} {
    ticker := time.NewTicker(interval)
    ch := make(chan struct{})
    go func() {
        defer ticker.Stop()
        defer close(ch)
        for {
            select {
            case <-ticker.C:
                select {
                case ch <- struct{}{}:
                default: // 丢弃多余的令牌
                }
            case <-ctx.Done():
                return
            }
        }
    }()
    return ch
}

7. 竞态检测

7.1 启用竞态检测

# 编译时启用
go build -race -o app .
go test -race ./...

# 运行时检测
./app

7.2 常见竞态示例

// 竞态:并发读写共享变量
var counter int

func main() {
    var wg sync.WaitGroup
    for i := 0; i < 1000; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            counter++ // 竞态!
        }()
    }
    wg.Wait()
    fmt.Println(counter) // 可能不是 1000
}

// 修复:使用原子操作
var counter int64
atomic.AddInt64(&counter, 1)
final := atomic.LoadInt64(&counter)

// 修复:使用互斥锁
var mu sync.Mutex
mu.Lock()
counter++
mu.Unlock()

// 修复:使用 channel
ch := make(chan int, 1000)
// 每个 goroutine: ch <- 1
// 汇总: for i := 0; i < 1000; i++ { <-ch }

7.3 原子操作

var count int64

// 加法
atomic.AddInt64(&count, 1)

// 读取
val := atomic.LoadInt64(&count)

// 写入
atomic.StoreInt64(&count, 100)

// 比较并交换(CAS)
swapped := atomic.CompareAndSwapInt64(&count, 100, 200)

// 自旋锁实现
type SpinLock struct {
    state int32
}

func (s *SpinLock) Lock() {
    for !atomic.CompareAndSwapInt32(&s.state, 0, 1) {
        runtime.Gosched() // 让出 CPU
    }
}

func (s *SpinLock) Unlock() {
    atomic.StoreInt32(&s.state, 0)
}

goroutine

基本写法:启动 goroutine go <函数>(<参数>)

// 启动 goroutine 执行函数
go doWork("task1");

基本写法:匿名函数 goroutine go func(<参数> <类型>) { ... }(<值>)

// 启动匿名函数 goroutine
go func(msg string) {
    fmt.Println(msg);
}("hello");

channel 创建

基本写法:无缓冲通道 make(chan <类型>)

// 无缓冲通道,发送和接收同步
ch := make(chan int);

基本写法:有缓冲通道 make(chan <类型>, <容量>)

// 有缓冲通道,容量为 10
ch := make(chan int, 10);

channel 操作

基本写法:发送数据 <通道> <- <值>

// 发送数据到通道
ch <- 42;

基本写法:接收数据 <-<通道>

// 从通道接收数据
v := <-ch;

基本写法:关闭通道 close(<通道>)

// 关闭通道,禁止再发送
close(ch);

基本写法:遍历通道 for <值> := range <通道> { ... }

// 遍历通道直到关闭
for v := range ch {
    fmt.Println(v);
}

select 语句

基本写法:select 多路复用 select { case ... }

// 多路复用选择
select {
case v := <-ch1:
    fmt.Println("ch1:", v);
case v := <-ch2:
    fmt.Println("ch2:", v);
case <-time.After(time.Second):
    fmt.Println("timeout");
}

基本写法:default 非阻塞 select { case ... default: }

// 非阻塞接收
select {
case v := <-ch:
    fmt.Println(v);
default:
    fmt.Println("no data");
}

sync.WaitGroup

基本写法:WaitGroup 等待 var <变量名> sync.WaitGroup

// 使用 WaitGroup 等待所有 goroutine 完成
var wg sync.WaitGroup;
for i := 0; i < 5; i++ {
    wg.Add(1);
    go func(n int) {
        defer wg.Done();
        doWork(n);
    }(i);
}
wg.Wait();

sync.Mutex

基本写法:互斥锁 var <变量名> sync.Mutex

// 互斥锁保护共享数据
var mu sync.Mutex;
var counter int;

func increment() {
    mu.Lock();
    defer mu.Unlock();
    counter++;
}

基本写法:读写锁 var <变量名> sync.RWMutex

// 读写锁,读多写少场景
var rwmu sync.RWMutex;
var data map[string]string;

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;
}

sync.Once

基本写法:单次执行 var <变量名> sync.Once

// sync.Once 确保初始化只执行一次
var (
    once sync.Once;
    instance *Config;
)

func GetConfig() *Config {
    once.Do(func() {
        instance = loadConfig();
    });
    return instance;
}

sync.Cond

基本写法:条件变量 sync.NewCond(&<互斥锁>)

// 条件变量等待通知
var mu sync.Mutex;
cond := sync.NewCond(&mu);

func waitForData() {
    mu.Lock();
    for !dataReady {
        cond.Wait();
    }
    mu.Unlock();
}

func notifyData() {
    mu.Lock();
    dataReady = true;
    cond.Signal();
    mu.Unlock();
}

sync.Pool

基本写法:对象池 sync.Pool{ New: func() any { ... } }

// 对象池复用对象
var bufPool = sync.Pool{
    New: func() any {
        return new(bytes.Buffer);
    },
}

func process(data []byte) {
    buf := bufPool.Get().(*bytes.Buffer);
    defer bufPool.Put(buf);
    buf.Reset();
    buf.Write(data);
}

atomic 原子操作

基本写法:原子加法 atomic.AddInt64(&<变量>, <值>)

// 原子加法
var counter int64;
atomic.AddInt64(&counter, 1);

基本写法:原子加载 atomic.LoadInt64(&<变量>)

// 原子读取
val := atomic.LoadInt64(&counter);

基本写法:原子存储 atomic.StoreInt64(&<变量>, <值>)

// 原子写入
atomic.StoreInt64(&counter, 100);

基本写法:原子比较交换 atomic.CompareAndSwapInt64(&<变量>, <旧值>, <新值>)

// CAS 操作
ok := atomic.CompareAndSwapInt64(&counter, 100, 200);

context

基本写法:创建根 context context.Background()

// 创建根 context
ctx := context.Background();

基本写法:带超时的 context context.WithTimeout(<父context>, <时长>)

// 创建 5 秒超时的 context
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second);
defer cancel();

基本写法:带取消的 context context.WithCancel(<父context>)

// 创建可取消的 context
ctx, cancel := context.WithCancel(context.Background());
defer cancel();

基本写法:带值的 context context.WithValue(<父context>, <键>, <值>)

// 创建带值的 context
ctx := context.WithValue(context.Background(), "userID", 123);

基本写法:从 context 获取值 <ctx>.Value(<键>)

// 从 context 获取值
userID := ctx.Value("userID").(int);

基本写法:检查 context 是否取消 <-<ctx>.Done()

// 检查 context 是否已取消
select {
case <-ctx.Done():
    return ctx.Err();
default:
    // 继续工作
}

并发模式

基本写法:fan-out 扇出 go <函数>(<输入通道>, <输出通道>)

// 多个 goroutine 处理同一输入
func worker(id int, jobs <-chan int, results chan<- int) {
    for j := range jobs {
        results <- j * 2;
    }
}

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

基本写法:fan-in 扇入 func <函数名>(<输出通道>, <输入通道1>, <输入通道2>)

// 合并多个通道的数据
func merge(out chan<- int, cs ...<-chan int) {
    var wg sync.WaitGroup;
    wg.Add(len(cs));
    for _, c := range cs {
        go func(ch <-chan int) {
            defer wg.Done();
            for v := range ch {
                out <- v;
            }
        }(c);
    }
    go func() {
        wg.Wait();
        close(out);
    }();
}

基本写法:pipeline 管道 go func() { ... }()

// 管道模式
func generate(nums ...int) <-chan int {
    out := make(chan int);
    go func() {
        defer close(out);
        for _, n := range nums {
            out <- n;
        }
    }();
    return out;
}

并发安全

基本写法:并发安全 map sync.Map

// 并发安全的 map
var m sync.Map;
m.Store("key", "value");
v, ok := m.Load("key");
m.Delete("key");

基本写法:遍历 sync.Map <map>.Range(func(<键>, <值>) bool { ... })

// 遍历 sync.Map
m.Range(func(key, value any) bool {
    fmt.Printf("%v: %v\n", key, value);
    return true;
});