Skip to content
Go back

Go 并发:通过通信共享内存

Updated:
Edit page

不要通过共享内存来通信;相反,通过通信来共享内存。

Do not communicate by sharing memory; instead, share memory by communicating.

这句话出自 Go 官方文档 Effective Go — Share by communicating。它强调了一种组织并发协作的方式:通过 channel 交接数据的处理权,减少多个 goroutine 同时修改共享状态。

例如,生产者生成任务,把它交给工作 goroutine;工作 goroutine 处理后,再把结果交给调用方。与其让各方反复检查共享变量里的“任务是否就绪”,不如直接用通信表达交接过程。

沿着这条主线,本文介绍三个相互配合的工具:channel 传递数据并同步执行,select 等待多个通信事件,context 传播取消和截止时间。至于共享状态是否需要锁,要根据具体问题选择,并不需要全部改成 channel。

一、从 goroutine 到 channel:执行与交接

go f() 会启动一个由 Go runtime 管理的 goroutine,让 f 与调用方并发执行。但启动并不等于等待完成:main 返回时,程序就会退出,不会自动等待其他 goroutine。

channel 可以让两个 goroutine 交换数据,并在交接处建立同步关系:

func main() {
    results := make(chan int)

    go func() {
        results <- 100 // 发送结果
    }()

    value := <-results // 等待并接收结果
    fmt.Println(value) // 100
}

make(chan int) 创建一个传递 int 的无缓冲 channel;ch <- value 发送,value := <-ch 接收。这里即使主 goroutine 先运行到接收处,也会等待结果,而不需要用 time.Sleep 猜测工作何时结束。

不过,收到一条消息,只能说明相应的交接已经发生。如果发送方在发送后还有其他工作,接收方并不能据此判断那些工作也已完成。

二、无缓冲与有缓冲:直接交接还是暂存

无缓冲 channel:双方配对

无缓冲 channel 没有暂存元素的空间,发送与接收必须配对。只有发送方或只有接收方到达时,先到的一方会等待。

sender -- value --> receiver
       rendezvous

因此,无缓冲 channel 适合表达直接交接。发送操作返回时,相应的接收已经发生,但接收方之后的业务处理可能才刚开始。

有缓冲 channel:允许有限积压

make(chan int, 3) 创建容量为 3 的 channel。在没有接收者时,前三次发送可以完成,第四次发送会等待空位:

ch := make(chan int, 3)
ch <- 10
ch <- 20
ch <- 30
// ch <- 40 // 此时会阻塞,直到有其他 goroutine 接收

fmt.Println(<-ch) // 10,腾出一个位置
ch <- 40

对于未关闭、非 nil 的有缓冲 channel,可以记住下面的规则:

操作条件行为
发送缓冲区未满可以完成,无需等待接收方
发送缓冲区已满等待接收方腾出位置
接收缓冲区有数据取出一个值
接收缓冲区为空等待发送方提供数据

容量用于吸收短暂的速度差异,并不能提高消费者本身的处理能力。下游持续较慢时,缓冲区最终会满,发送方就会被阻塞,从而形成背压(下游处理不及时,通过阻塞等方式限制上游生产速度,避免任务持续积压)。

这个约束要真正传回生产流程才有效。如果每生成一个任务就启动一个新的 goroutine 去发送,积压可能转移到大量等待发送的 goroutine 上,整体资源使用仍然没有上限。

三、关闭 channel:声明发送结束

close(ch) 表示以后不会再向这个 channel 发送数据。它不会丢弃已有的缓冲数据:接收方仍可读完它们,之后接收才会得到元素类型的零值和 ok == false

ch := make(chan int, 2)
ch <- 10
close(ch)

value, ok := <-ch // 10, true
fmt.Println(value, ok)
value, ok = <-ch  // 0, false
fmt.Println(value, ok)

因此,不能仅凭读到 0 判断 channel 是否关闭。处理一串数据时,for range 会持续接收,直到 channel 关闭且数据读完:

func producer(out chan<- int) {
    defer close(out)
    for i := 0; i < 5; i++ {
        out <- i
    }
}

func consumer(in <-chan int) {
    for value := range in {
        fmt.Println(value)
    }
}

这里 chan<- int 只允许发送,<-chan int 只允许接收,用类型标明了职责。调用方需要让生产与消费并发运行,例如先 go producer(ch),再执行 consumer(ch)

关闭应由能够确认“所有发送都已结束”的一方负责。 单个生产者通常自己关闭;多个生产者共用一个 channel 时,应由协调方等待所有生产者结束后统一关闭,不能让每个生产者各自 close

几个边界可以一起记住:向已关闭的 channel 发送、重复关闭、关闭 nil channel 都会 panic;对 nil channel 发送或接收则会一直阻塞。channel 也不要求用完必须关闭,只有接收方需要“发送结束”这个信号时才需要关闭。

四、channel 也是同步机制

有时接收方不需要业务数据,只需要知道某项工作已经结束。这时可以用 chan struct{} 表达完成信号:

var result int
done := make(chan struct{})

go func() {
    result = 100
    close(done)
}()

<-done
fmt.Println(result) // 100

这个例子里的 result 虽然是共享变量,但读取发生在完成信号之后。Go 内存模型规定,关闭 channel 与因该关闭而返回的接收之间存在同步关系,所以关闭前的写入对接收后的读取可见。这种先后关系称为 happens-before

发送与对应的接收也能建立同步关系。因此,channel 的作用不只是传递一个值,还能让双方对“哪些操作已经完成”达成一致。

但这个保证有明确边界:如果发送的是指针、切片等引用同一份数据的值,channel 不会自动禁止发送方继续修改底层数据。双方仍需约定处理权如何交接,避免交接后又并发读写。

五、select:给等待增加其他出口

单独执行 <-ch 只能等待这个 channel。select 可以同时等待多个发送或接收操作,例如同时等待结果和超时:

select {
case value, ok := <-results:
    if !ok {
        return // 结果流已经结束
    }
    fmt.Println(value)
case <-time.After(time.Second):
    fmt.Println("等待结果超时")
}

按照 Go 语言规范,如果多个通信分支都可以执行,select 会从中均匀伪随机选择一个,没有源码顺序上的优先级。如果都不能执行,就阻塞等待;如果有 default,则转而执行它。

例如,下面的代码只尝试接收一次,不会等待:

select {
case value, ok := <-results:
    fmt.Println(value, ok)
default:
    fmt.Println("当前没有可接收的结果")
}

不要把这种写法直接放进没有其他阻塞操作的无限循环,否则循环可能不断执行 default,形成空转。已关闭的 channel 则始终可以接收;在循环里处理多个 channel 时,应在发现关闭后退出,或将对应 channel 变量设为 nil,使该分支不再参与选择。

还要区分:上面的超时只让当前等待结束,并不会自动停止后台工作。要把停止请求传给工作方,需要双方约定取消机制。

六、context:把取消与超时传下去

context.Context 用于沿调用链传播取消、截止时间和请求范围的数据。它的 Done() 返回一个 channel:当 context 被取消或到达截止时间时,这个 channel 会关闭,等待者就能收到通知。

把普通发送改为同时监听取消,可以让发送方在下游停止接收时有机会退出:

func send(ctx context.Context, out chan<- string, value string) error {
    select {
    case out <- value:
        return nil
    case <-ctx.Done():
        return ctx.Err()
    }
}

调用方通过 WithTimeout 为这次等待设置期限。下面是完整示例;因为没有接收方,发送会在超时后返回:

package main

import (
    "context"
    "fmt"
    "time"
)

func send(ctx context.Context, out chan<- string, value string) error {
    select {
    case out <- value:
        return nil
    case <-ctx.Done():
        return ctx.Err()
    }
}

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
    defer cancel()

    out := make(chan string)
    fmt.Println(send(ctx, out, "report")) // context deadline exceeded
}

创建派生 context 后,通常立即用 defer cancel() 安排退出时释放资源;如果工作提前结束,也可以更早调用 cancel()WithCancel 用于主动取消,WithTimeout 设置相对时长,WithDeadline 设置具体截止时间。父 context 取消时,派生出的子 context 也会被取消。context 官方文档说明了这些传播规则。

取消是协作式通知,不会强制终止 goroutine。 工作方必须检查 Done()Err(),或调用支持 context 的 API;调用方若要确认它已经退出,还需要等待完成信号或使用 WaitGroup

取消分支也没有特殊优先级:如果 out 可以发送且 ctx.Done() 已关闭,上面的 select 仍可能选择发送。因此,这个模式提供取消出口,但不保证取消发生后绝不再发送数据。

传参时,把 ctx 放在函数的第一个参数,不传 nil,也通常不把它长期保存在结构体里。业务任务和结果用参数或 channel 传递,context 的值只用于请求范围的元数据。

七、什么时候用 channel,什么时候用锁

“通过通信共享内存”是一种设计思路,具体工具仍应围绕问题选择。Go 官方 Wiki也建议优先选择表达最清晰、最简单的方式。

需要解决的问题常用工具典型场景
交接数据或任务channel任务队列、结果流、完成信号
等待多个通信事件select同时等待结果、取消或超时
传播工作生命周期context请求取消、跨调用链超时
保护共享状态Mutex / RWMutex多个 goroutine 访问同一个 map
等待一组任务结束WaitGroup等全部 worker 退出
更新简单的独立状态atomic原子计数器、状态位

例如,Worker Pool 可以用 channel 分发任务,用 WaitGroup 等待全部 worker 退出,再由协调方关闭结果 channel;如果 worker 还要共同更新一张统计表,就用锁保护那张表。这些工具可以自然组合。

原子操作则只保证相应操作的原子性,多个字段之间的约束或一组操作的整体一致性,不能仅靠把每个字段都换成 atomic 来保证。更复杂的条件等待还可以使用 sync.Cond,本文不展开。

八、总结

理解 Go 的并发通信,可以从一次任务交接出发:goroutine 执行任务,channel 传递数据并建立同步关系,缓冲区容纳有限积压,关闭表示发送结束;当等待需要兼顾其他事件时,用 select 组织分支,用 context 传播取消和超时。“通过通信共享内存”的关键,是把数据由谁处理、何时交接、如何结束表达清楚。需要保护共享状态时使用锁,需要确认一组任务已经退出时使用 WaitGroup,让每种工具承担清晰的职责。

参考资料


Edit page