Golang 并发实践与模式
12.5 并发实践与模式(Pattern)
12.5.1 确保 API 不涉及并发
并发是一种实现细节,良好的 API 设计应该尽可能隐藏实现细节。这样可以在不改变代码调用方式的情况下,改变代码的内部工作方式。
实际上,这意味着永远不应该在 API 的类型、函数和方法中暴露通道或互斥锁(将在“12.6 用互斥锁(mutex)?还是通道(channel)?<274页>”一节介绍互斥锁)。如果暴露了通道,就会将通道管理的责任交给了 API 的使用者。API 使用者就必须担心诸如通道是否有缓冲、通道是否已关闭或通道是否为 nil 等问题。使用者还可能通过以意外的顺序访问通道或互斥锁(mutex)来引发死锁。
这并不意味着你永远不能把通道(channel)作为函数参数或结构体字段,而是说通道不应该被导出(exported)。
这条规则(不暴露通道)也存在一些例外情况。如果 API 是一个包含并发辅助函数的库,那么通道将成为其 API 的一部分。
12.5.2 协程、for 循环与可变变量
用于启动协程的闭包(closures)通常没有参数,而是通过捕获其声明环境中的变量(如引用外部作用域的变量)来工作。在 Go 1.22 之前,当闭包尝试捕获 for 循环的索引或值时,会出现意外行为。因为 for 循环每次迭代会重复使用同一个索引/值变量,闭包捕获的是同一个变量,导致所有协程引用的是循环结束时的最终值,而非期望的当前迭代值。如“4.4.5.3 for-range 得到的值是副本<71页>”一节所述,Go 1.22 改变了 for 循环的行为,每次迭代会创建新的索引/值变量(即变量隔离)。闭包捕获的是当前迭代的独立变量,避免了多协程引用同一变量的问题,提升了代码的预期行为。
下面这段代码演示了为什么 Go 1.22 的这个改进是值得的。
如果在 Go 1.21 或更早的版本上运行下面的代码,或者在 Go 1.22 或更高版本上运行,但通过在 go.mod 文件的 go 指令中将 Go 版本设置为 1.21 或更早来模拟旧版本的行为,就会发现一个不易察觉的 bug。
func main() {
a := []int{2, 4, 6, 8, 10}
ch := make(chan int, len(a))
for _, v := range a {
go func() {
ch <- v * 2
}()
}
for i := 0; i < len(a); i++ {
fmt.Println(<-ch)
}
}
为数组 a 中的每个值启动一个协程。看起来像是为每个协程传递了不同的值,但运行代码后会发现实际情况并非如此。
20
20
20
20
20
Go 1.21 及更早版本中,所有协程向通道 ch 写入的值都是 20,是因为 for 循环的索引和值变量是共享的,闭包捕获的是引用(对同一个变量的引用),而非当前迭代的独立值。最终 for 循环结束时,捕获的 v 已经是循环结束时的值 10(即数组 a 的最后一个元素)。
升级到 Go 1.22 或更高版本并将 go.mod 中 go 指令的值更改为 1.22 或更高版本会更改 for 循环的行为,以便在每次迭代时创建新的索引和值变量。这将为每个 goroutine 传递不同的值,使程序运行结果符合预期:
20
8
4
12
16
如果无法升级到 Go 1.22,可以通过两种方法解决闭包捕获循环变量的错误:
- 在循环内部通过变量遮蔽(Variable Shadowing)创建值的副本。
for _, v := range a {
v := v // 创建值的副本,而不是捕获循环变量。
go func() {
ch <- v * 2
}()
}
- 将值作为参数传递给协程。此方法可避免使用变量遮蔽(Variable Shadowing),让数据流向更清晰明了。
for _, v := range a {
go func(val int) {
ch <- val * 2
}(v)
}
虽然 Go 1.22 解决了 for 循环中索引和值变量被闭包捕获的问题,但你仍然需要小心其他被闭包捕获的变量。无论闭包是否用于协程,只要它依赖的变量值可能变化,就必须通过参数传递该值给闭包,或者确保每个引用该变量的闭包都有自己的独立副本。
任何时候闭包(closures)使用的变量值可能发生变化时,都应该通过参数将变量的当前值传递给闭包,以确保闭包使用的是该变量的即时副本,而不是依赖可能被后续修改的外部变量。
12.5.3 务必妥善处理协程退出
每次启动一个协程函数时,都必须确保它最终会正常退出。与变量不同,Go 运行时无法检测到某个协程是否不再被使用。如果协程没有正常退出,其栈上分配的所有变量内存将保持分配状态,且任何由该协程栈变量引用的堆内存都无法被垃圾回收。这种情况被称为“协程泄漏(goroutine leak)”。
协程是否会正常退出,并不总是显而易见的。例如,假设使用协程作为生成器(不断产生数据的协程):
func countTo(max int) <-chan int {
// 无缓冲通道
ch := make(chan int)
go func() {
for i := 0; i < max; i++ {
ch <- i
}
close(ch)
}()
return ch
}
func main() {
for i := range countTo(10) {
fmt.Println(i)
}
}
这个例子只是一个简短的演示;在实际编程中,不要用协程去生成一个数字列表。因为生成数字列表的操作本身非常简单,直接在主程序里完成就足够了。用协程去生成数字列表,违背了之前提到的一条关于“何时使用并发”的指导原则。
在通常情况下(即你使用所有生成的值时),协程会正常退出。然而,如果提前退出循环,协程就会永远阻塞,因为它在等待有人从通道读取值,但此时已经没有读取方了。
func main() {
for i := range countTo(10) {
if i > 5 {
break
}
fmt.Println(i)
}
}
12.5.4 用 Context 终止协程
要解决 countTo 协程泄漏的问题,需要找到一种方法来告诉协程何时该停止处理。在 Go 语言中,通常使用 context(上下文)机制来实现这个目的。下面这段重写的 countTo 代码就是用来演示这种技术的,可以在本章资源库 sample_code/context_cancel 目录下找到这段代码。
func countTo(ctx context.Context, max int) <-chan int { // ❶ 增加 context.Context 参数。
ch := make(chan int)
go func() {
defer close(ch)
for i := 0; i < max; i++ {
select {
case <-ctx.Done(): // ❷ 检查是否应该退出协程。如果 Done 通道返回了值,就退出循环。
return
case ch <- i: // ❸ 向通道 ch 写入数据
}
}
}()
return ch
}
func main() {
ctx, cancel := context.WithCancel(context.Background()) // ❹ 创建 context 和 cancel 函数。
defer cancel() // ❺ 确保在 main 函数退出时调用 cancel。
ch := countTo(ctx, 10)
for i := range ch {
if i > 5 {
break
}
fmt.Println(i)
}
}
countTo 函数被修改为除了 max 参数外,还接收一个 context.Context 参数(❶处)。协程中的 for 循环调整为 for-select 循环,包含两个 case 分支。其中一个分支尝试向 ch 写入数据(❸处),另一个分支检查由 context.Done 方法返回的通道(❷处)。如果该通道返回了值,就退出 for-select 循环以及协程。现在,就有了一种方法,可以在所有值都被读取完毕时防止协程泄漏。
这就引出了一个新问题:如何让 Done 通道返回值?这是通过上下文取消(Context Cancellation)触发的。在 main 函数中,通过调用 context.WithCancel 函数(❹处)来创建一个 context 和一个 cancel 函数。接着,使用 defer 确保在 main 函数退出时会调用 cancel 函数(❺处)。这会关闭由 Done 返回的通道,而由于已关闭的通道总是会返回值,这就确保了运行 countTo 的协程能够正常退出。
使用 context 来终止协程是一种很常见的编程模式。它允许开发者根据调用栈中更早的函数发出的信号来停止协程。在“14.3 上下文取消(Context Cancellation)<314页>”一节中,将详细介绍如何使用 context 来告知一个或多个协程,是时候关闭了。
12.5.5 何时使用缓冲通道与无缓冲通道
在 Go 的并发编程中,最难掌握的技术之一就是判断何时使用缓冲通道。默认情况下,通道是无缓冲的,它们易理解:一个 goroutine 写入数据并等待另一个 goroutine 读取,就像接力赛中的接力棒一样。而缓冲通道复杂得多,因为必须指定缓冲区大小,而缓冲通道的缓冲区大小是有限的。正确使用缓冲通道意味着必须处理缓冲区已满的情况,此时写入的 goroutine 会阻塞等待读取的 goroutine。那么,如何正确使用缓冲通道呢?
缓冲通道的适用场景微妙,总结为:当确切知道已启动的 goroutine 数量、想限制将要启动的 goroutine 数量或想限制排队的工作量时,缓冲通道是有用的。
缓冲通道在以下两种情况下表现良好:
- 想从已启动的一组 goroutine 中收集数据。
- 想限制并发使用。
缓冲通道还有助于管理系统中排队的工作量,防止服务落后并过载。这里会给出几个例子来说明缓冲通道的使用方式。
第一个示例,需要处理通道中的前 10 个结果。为此,启动 10 个 goroutine,每个 goroutine 负责处理一个任务并将处理后的结果写入同一个缓冲通道(buffered channel)。
func processChannel(ch chan int) []int {
const conc = 10
results := make(chan int, conc) // ❶ 创建一个大小为 10 的缓冲通道。
for i := 0; i < conc; i++ { // ❷ 启动 10 个 goroutine。
go func() {
v := <-ch
results <- process(v) // ❸ 每个 goroutine 都将 process 结果写入缓冲通道。
}()
}
var out []int
for i := 0; i < conc; i++ { // ❹ 循环读取缓冲通道 results 中的值。
out = append(out, <-results)
}
return out
}
在此示例(完整代码见本章资源库 sample_code/buffered_channel_work 目录)中,已确切知道启动了 10 个 goroutine,并且希望每个 goroutine 在完成工作后立即退出。这意味着可以创建一个大小为 10(与协程数量相同)的缓冲通道(❶处),其缓冲区大小为每个启动的 goroutine 提供了空间,并让每个 goroutine 向此缓冲通道写入数据而不会被阻塞(❸处)。然后,可以循环遍历这个缓冲通道(❹处),读取被写入的值。当所有值都被读取后,则返回结果,并知道没有造成任何 goroutine 泄漏。
12.5.6 用缓冲通道实现背压(Backpressure)
使用缓冲通道可以实现的另一种技术是背压(Backpressure)。背压(Backpressure)机制通过限制组件的工作量来提高系统整体性能,尽管这看起来有些反直觉。可以使用一个缓冲通道和一个 select 语句来限制系统中同时请求的数量。
// 实现了一个令牌桶限流机制,通过缓冲通道控制并发访问量。
type PressureGauge struct {
ch chan struct{} // ch 是一个缓冲通道,struct{} 作为令牌(不占内存)
}
// New 创建并返回一个新的 PressureGauge 实例,limit 参数指定最大并发处理能力。
func New(limit int) *PressureGauge {
return &PressureGauge{
ch: make(chan struct{}, limit), // 容量由 limit 决定,代表最大并发数
}
}
// Process 用于执行一个函数 f,但会先检查当前是否有足够的并发处理能力。
// 如果有足够的能力,则执行 f 并返回 nil;否则返回错误。
func (pg *PressureGauge) Process(f func()) error {
select {
case pg.ch <- struct{}{}: // 向通道写入一个空结构体(令牌),表示占用一个处理能力
f() // 执行传入的函数
<-pg.ch // 从通道接收一个空结构体(令牌),表示释放一个处理能力
return nil // 返回 nil 表示成功执行
default: // 如果通道已满(即达到最大并发限制),则执行 default 分支
return errors.New("no more capacity") // 返回错误,表示没有更多处理能力
}
}
PressureGauge是有状态的,不是因为struct包含chan,而是因为Process()会改变ch所指向 channel 的占用状态,并且下一次Process()是否成功取决于这个当前状态。但可以构造出“struct 有 chan 字段,但这个类型本身不承担状态语义”的例子:
type Config struct { Notify chan<- Event } func NewService(cfg Config) *Service { return &Service{ notify: cfg.Notify, } }这里
Config虽然含有chan,但Config本身只是描述:Service 应该往哪个 channel 发通知。一般不会因此把Config当作一个“有状态对象”。再比如:
type Dependencies struct { Events chan<- Event Logs chan<- Log }它可能纯粹是依赖注入容器。channel 自身有状态,但
Dependencies的业务抽象并不是“维护状态”。所以最核心的判断依据还是看数据类型的行为(即方法),而不是看数据类型是什么。
根据校验,通常情况下,如果
chan/slice/map等引用类型是这个对象内部运行机制的一部分,通常认为对象有状态;如果 ``chan/slice/map只是一个外部依赖、配置或句柄,则对象不是有状态的。
这段代码实现了一个令牌桶限流机制,通过缓冲通道控制并发访问量。代码中,创建了一个结构体,其中包含一个可以容纳一定数量“令牌”(用空结构体表示)的缓冲通道和一个待执行的函数 f。每当一个 goroutine 想要使用这个函数时,它会调用 Process 方法。这是一个罕见的例子——同一个 goroutine 既要读取又要写入同一个通道。select 尝试向通道写入一个令牌:
- 如果写入成功,函数
f就会被执行;然后,从缓冲通道中读取一个令牌(表示归还)。 - 如果写入失败(即通道已满),就会执行
default分支,并返回一个错误。
下面有一个快速示例(完整代码见本章资源库 sample_code/backpressure 目录),展示了如何将这段代码与内置的 HTTP 服务器一起使用(将在“13.4.2 服务器<295页>”一节介绍更多关于处理 HTTP 的知识)。
func doThingThatShouldBeLimited() string {
time.Sleep(2 * time.Second)
return "done"
}
func main() {
pg := New(10)
http.HandleFunc("/request", func(w http.ResponseWriter, r *http.Request) {
err := pg.Process(func() {
w.Write([]byte(doThingThatShouldBeLimited()))
})
if err != nil {
w.WriteHeader(http.StatusTooManyRequests)
w.Write([]byte("Too many requests"))
}
})
http.ListenAndServe(":8080", nil)
}
12.5.7 关闭 select 中的 case 分支
select 非常适合从多个并发通道读取数据,但需注意已关闭通道的处理。如果 select 的某个 case 分支读取一个已关闭的通道,该分支会持续成功并返回零值,导致程序总是误处理无效数据。因此,每次该分支被选中时,都需要检查值的有效性,并跳过该分支。如果读取操作是分散进行的,程序将会浪费大量时间读取垃圾值。即使其他通道活跃,程序仍然会耗费时间从已关闭的通道读取数据,因为 select 是随机选择每个 case 分支的。
读取 nil 通道会导致代码永久阻塞(类似错误),但可利用此特性禁用 select 分支。检测到通道关闭后,将其变量设为 nil,关联的分支将不再执行,因为 nil 通道的读取永远不会返回值。如下示例(完整代码见 The Go Playground 或本章资源库 sample_code/close_case)展示了一个 for-select 循环,从两个通道读取数据,直到它们都关闭。通过将已关闭通道设为 nil,确保程序仅处理有效数据。
nil 通道永远不会被 case 选择,因为 nil 通道不可读写
// in and in2 are channels
for count := 0; count < 2; {
select {
case v, ok := <-in:
if !ok {
in = nil // the case will never succeed again!
count++
continue
}
// process the v that was read from in
case v, ok := <-in2:
if !ok {
in2 = nil // the case will never succeed again!
count++
continue
}
// process the v that was read from in2
}
}
12.5.8 超时控制
大多数交互式应用程序都必须在特定时间内返回响应,否则会影响用户体验。Go 通过并发机制管理请求执行时间。Java/JavaScript 等语言需要额外引入“超时回调”等机制实现超时逻辑,而 Go 直接利用 select + 通道的现有特性即可构建超时逻辑,无需额外库。如下所示(完整示例代码见 The Go Playground 或本章资源库 sample_code/time_out 目录):
func timeLimit[T any](worker func() T, limit time.Duration) (T, error) {
out := make(chan T, 1)
ctx, cancel := context.WithTimeout(context.Background(), limit)
defer cancel()
go func() {
out <- worker()
}()
select {
case result := <-out: // ❶ 等待 worker 完成,并从 out 通道读取结果。
return result, nil
case <-ctx.Done(): // ❷ 如果 Done 通道返回一个值,则表示超时。
var zero T
return zero, errors.New("work timed out")
}
}
Go 中限制操作时间时,通常会使用基于 select 和 context 的模式。这种模式是 Go 并发编程中的常见惯用法,用于实现超时控制。“十四 上下文(Context)<305页>”中将详细介绍上下文,并展示如何在超时情况下优雅地终止 goroutine。本章只需要了解,当操作达到超时时间时,上下文会被自动取消。context.Done() 返回一个通道,当上下文因超时或显式调用上下文 cancel 方法而取消时,该通道会返回一个值。可以使用 context.WithTimeout 函数创建一个带超时的上下文,并使用 time 包(详见“13.2 time 包<284页>”一节)中的常量来指定等待时间。
一旦设置好上下文,就可将耗时操作(worker)放到一个 goroutine 中执行,避免阻塞主线程。然后,主线程使用 select 在两个 case 分支之间进行选择:
- 第一个
case分支(❶处)是在工作(worker)完成时从out通道读取值。 - 第二个
case分支(❷)是等待Done方法返回的通道返回一个值,就像在“12.5.4 用 Context 终止协程<262页>”中看到的那样。如果Done返回的通道确实返回了值,就返回一个超时错误。
写入一个大小为 1 的缓冲通道,这样即使 Done 分支首先被触发,goroutine 中的通道写入也能完成。如果 timeLimit 在 goroutine 处理完成之前就退出了,那么 goroutine 会继续运行,最终会将返回的值写入缓冲通道并退出。你只是不对返回的结果做任何处理。如果想在不再等待 goroutine 完成时停止它的工作,请使用上下文取消(Context Cancellation),详见“14.3 上下文取消(Context Cancellation)<314页>”一节。
12.5.9 使用 WaitGroups
有时候,一个 goroutine 需要等待多个 goroutine 完成它们的工作。如果等待单个 goroutine,可以使用之前介绍的上下文取消(Context Cancellation)模式。但是,如果需要等待多个 goroutine 完成,就需要使用标准库中的 sync.WaitGroup。下面是一个简单示例(完整示例代码见 The Go Playground 或本章资源库 sample_code/waitgroup 目录)。
func main() {
var wg sync.WaitGroup
wg.Add(3)
go func() {
defer wg.Done()
doThing1()
}()
go func() {
defer wg.Done()
doThing2()
}()
go func() {
defer wg.Done()
doThing3()
}()
wg.Wait()
}
sync.WaitGroup 不需要初始化,只需要声明即可,因为它的零值(Zero Value)是有用的。sync.WaitGroup 上有三个方法:
Add,用于增加等待的 goroutine 计数器;Add通常只调用一次,传入将要启动的 goroutine 数量。Done,用于减少计数器,并在 goroutine 完成时被调用;Done在 goroutine 内部调用。为了确保即使 goroutine 发生 panic 也会调用Done,需要用defer调用Done方法。Wait,它会暂停其 goroutine,直到计数器归零。
如您所见,代码中没有显式地传递 sync.WaitGroup。这有两个原因:
- 第一个原因是,必须确保所有使用
sync.WaitGroup的地方都使用同一个实例。如果将sync.WaitGroup传递给 goroutine 函数,并且不使用指针,那么函数会得到一个sync.WaitGroup副本,对Done的调用将不会减少原始的sync.WaitGroup。通过使用闭包(closures)来捕获sync.WaitGroup,可以确保每个 goroutine 都引用的是同一个sync.WaitGroup实例。 - 第二个原因是设计上的考虑。请记住,应该将并发逻辑排除在 API 之外。就像之前在通道中看到的那样,通常的模式是使用一个闭包(closures)来启动 goroutine,这个闭包封装了业务逻辑。闭包处理与并发相关的问题,而函数则提供算法。
让我们来看一个更现实的例子。如前所述,当有多个 goroutine 向同一个通道写入数据时,需要确保该通道只被关闭一次。sync.WaitGroup 非常适合用于这种场景。让我们看看它在下面这个函数中是如何工作的:该函数会并发地处理通道中的值,将结果收集到一个切片中,并返回该切片。
func processAndGather[T, R any](in <-chan T, processor func(T) R, num int) []R {
out := make(chan R, num)
var wg sync.WaitGroup
wg.Add(num)
for i := 0; i < num; i++ {
go func() {
defer wg.Done()
for v := range in {
out <- processor(v)
}
}()
}
go func() { // ❶ 监控 goroutine,等待所有处理 goroutine 退出。
wg.Wait()
close(out)
}()
var result []R
for v := range out { // ❷
result = append(result, v)
}
return result
}
在此示例(完整代码见本章资源库 sample_code/waitgroup_close_once 目录)中,启动了一个监控 goroutine,它会等待所有处理 goroutine 退出(❶处)。当它们都退出后,监控 goroutine 会调用 close 来关闭输出通道。当 out 被关闭且缓冲区为空时,for-range 通道循环就会退出(❷处)。最后,该函数返回处理过的值。
虽然 WaitGroup 很方便,但它们不应该是协调 goroutine 的首选。只有在所有工作 goroutine 退出后需要清理某些资源(比如关闭工作 goroutine 都写入的同一个通道)时,才应该使用 WaitGroup。
golang.org/x 与 errgroup.Group
Go 的开发者维护了一套补充标准库的工具集。这套工具集统称为
golang.org/x包,其中包含一个errgroup包,该包内有一个errgroup.Group类型。它建立在WaitGroup之上,用于创建一组 goroutine,当其中任何一个 goroutine 返回错误时,这些 goroutine 就会停止处理。请阅读errgroup.Group文档以了解更多信息。
12.5.10 确保代码精确执行一次
正如“10.3.10 尽量避免使用 init 函数<209页>”一节所述,init 函数应当仅用于初始化那些实质上不可变(immutable)的包级变量,但有时需要延迟初始化,或者在程序启动后恰好只执行一次初始化代码。这通常是因为初始化过程相对耗时,而且程序可能并不是每次运行都需要执行初始化。sync.Once 是 Go 标准库提供的一个类型,可确保某段代码在并发环境下仅执行一次。让我们快速看一下它是如何工作的。假设有一段初始化比较耗时的代码:
type SlowComplicatedParser interface {
Parse(string) string
}
func initParser() SlowComplicatedParser {
// do all sorts of setup and loading here
}
下面是如何使用 sync.Once 来延迟初始化一个 SlowComplicatedParser(慢且复杂的解析器)的示例。
var parser SlowComplicatedParser
var once sync.Once
func Parse(dataToParse string) string {
once.Do(func() {
parser = initParser()
})
return parser.Parse(dataToParse)
}
示例代码中定义了两个包级变量:parser(类型为 SlowComplicatedParser)和 once(类型为 sync.Once)。sync.Once 无需显式初始化,直接使用零值即可工作,这是 Go 语言中“零值可用”设计模式的体现(类似 sync.WaitGroup)。
与 sync.WaitGroup 一样,必须确保不要复制 sync.Once 的实例,因为复制 sync.Once 会导致每个副本独立维护状态,无法保证初始化只执行一次。在函数内部声明 sync.Once 是错误的做法,因为每次调用都会创建新实例,无法记忆之前的执行情况。
sync.WaitGroup的定义:type WaitGroup struct { noCopy noCopy // Bits (high to low): // bits[0:32] counter // bits[32] flag: synctest bubble membership // bits[33:64] wait count state atomic.Uint64 sema uint32 }A WaitGroup must not be copied after first use.
在示例中,希望确保 parser 只被初始化一次,因此通过 once.Do 方法传入闭包,在闭包中完成 parser 的初始化。无论 Parse 被调用多少次,once.Do 只会执行闭包一次,确保 parser 只初始化一次。
Go 1.21 新增了三个辅助函数(sync.OnceFunc、sync.OnceValue 与 sync.OnceValues),使得函数恰好运行一次变得更加容易。它们的唯一区别在于传入函数的返回值数量(0 个、1 个或 2 个)。sync.OnceValue 与 sync.OnceValues 是泛型(Generics)函数,能自动适应原始函数的返回值类型(如 int、string、自定义类型等),无需手动指定类型。
sync 包的这些函数使用方法非常简单。将原始函数传递给辅助函数,然后会得到一个新函数,该辅助函数只会调用一次原始函数。原始函数返回的值会被缓存。下面展示如何用 sync.OnceValue 重写前文示例中的 Parse 函数(完整代码见 The Go Playground 或本章资源库 sample_code/sync_value 目录)。
var initParserCached func() SlowComplicatedParser = sync.OnceValue(initParser)
func Parse(dataToParse string) string {
parser := initParserCached()
return parser.Parse(dataToParse)
}
initParserCached 是一个包级变量,它存储了 sync.OnceValue(initParser) 返回的函数。首次调用 initParserCached 函数时,initParser 函数会被执行,并且 initParser 的返回值会被缓存。后续再调用 initParserCached 时,都直接返回缓存的值,不再重复执行 initParser 函数。这意味着可以移除包级变量 parser,因为通过 sync.OnceValue,初始化逻辑和返回值的管理完全由标准库处理,不再依赖包级变量。
12.6 用互斥锁(mutex)?还是通道(channel)?
其他编程语言在协调多个线程对数据的访问时,可能要用到互斥锁(mutex)。互斥锁(mutex)是 “mutual exclusion” 的缩写,其作用是限制某些代码的并发执行,或者限制对共享数据的访问。这段受保护的部分被称为“临界区”。
Go 语言的创造者设计通道和 select 来管理并发是有充分理由的。互斥锁(mutex)的主要问题是它会掩盖程序中数据的流向。当值通过一系列通道在 goroutine 之间传递时,数据流向是清晰的。通道中的值一次只能被一个 goroutine 访问。当使用互斥锁(mutex)保护一个值时,无法确定该值当前被哪个 goroutine 拥有,因为所有并发进程都可以共享该值的访问权限。这使得理解数据的处理顺序变得困难。Go 社区有一句谚语来描述这种理念:“可以通过通信来共享内存,但不要通过共享内存来通信。”
标准库中的 sync 包提供了两种互斥锁的实现:
- 第一种是
Mutex,它有两个方法:Lock和Unlock。当调用Lock时,如果此时已经有另一个 goroutine 正在执行被这个锁保护的“临界区(critical section)”(即需要互斥访问的代码块或数据),那么当前这个调用Lock的 goroutine 就会暂停/阻塞,直到那个在临界区的 goroutine 执行完毕并调用了Unlock。一旦临界区空闲,锁就被当前 goroutine 获取。然后,当前 goroutine 就可以安全地执行临界区内的代码了。调用Unlock表示 goroutine 完成了对临界区的操作,锁被释放,其他等待的 goroutine 可以尝试获取锁,这标志着临界区的结束。 - 第二种是
RWMutex,它允许区分读锁(独占)与写锁(共享)。任何时刻只能有一个写者进入临界区(通过Lock获取写锁,Unlock释放写锁),其他写者必须等待。但是,多个读者可以同时持有读锁(通过RLock获取读锁,RUnlock释放读锁)。
必须确保每次获取锁后都能正确释放锁。使用 defer 语句在调用 Lock 或 RLock 后立即调用 Unlock,可确保互斥锁的正确释放:
type MutexScoreboardManager struct {
l sync.RWMutex
scoreboard map[string]int
}
func NewMutexScoreboardManager() *MutexScoreboardManager {
return &MutexScoreboardManager{
scoreboard: map[string]int{},
}
}
func (msm *MutexScoreboardManager) Update(name string, val int) {
msm.l.Lock()
defer msm.l.Unlock()
msm.scoreboard[name] = val
}
func (msm *MutexScoreboardManager) Read(name string) (int, bool) {
msm.l.RLock()
defer msm.l.RUnlock()
val, ok := msm.scoreboard[name]
return val, ok
}
现在,已经介绍了使用互斥锁实现的版本,在决定使用它们之前,请谨慎选择。Katherine Cox-Buday 的著作《Concurrency in Go(O’Reilly)》中包含了一个决策树,可帮助决定是使用通道,还是互斥锁:
- 如果正在协调 goroutine 或跟踪一个值在一系列 goroutine 中如何被转换,请使用通道。
- 如果正在共享对一个结构体中字段的访问,请使用互斥锁。
- 如果在使用通道时发现了关键的性能问题(详见“15.6 基准测试(Benchmark)<348页>”),并且找不到其他修复问题的方法,那么重构代码以使用互斥锁。
由于记分板(scoreboard)是结构体中的一个字段,且不会被传递给其他地方,这种情况下使用互斥锁是合理的,因为数据是内存中的。当数据存储在外部服务(如 HTTP 服务器或数据库)中时,不应使用互斥锁来保护访问。
互斥锁需要更多的手动管理(Bookkeeping)。例如,必须确保每个 Lock() 都对应一个 Unlock(),否则程序可能死锁。文中例子在同一个方法里完成加锁和解锁,避免跨方法调用时忘记释放锁。另一个问题是 Go 的互斥锁不支持“重入(Reentrant)”,即一个 goroutine 不能重复获取同一个锁。如果一个 goroutine 尝试获取同一个锁两次,则会导致死锁,因为它在等待自己释放锁。这与 Java 等语言不同,后者中的锁是可重入的(Reentrant),同一个线程可以多次获取同一把锁而不会死锁。
不可重入锁(Nonreentrant)使得在调用自身递归的函数中获取锁变得棘手。必须在递归函数调用之前释放锁。一般来说,在持有锁时进行函数调用要小心,因为不知道这些调用中会获取哪些锁。如果函数调用了另一个试图获取相同互斥锁的函数,该 goroutine 就会死锁。
与 sync.WaitGroup 和 sync.Once 一样,互斥锁(mutex)绝不能被复制。如果互斥锁被复制,它的锁将不会被共享。如果将互斥锁传递给一个函数或作为结构体中的一个字段被访问,必须通过指针进行。
Do not communicate by sharing memory; instead, share memory by communicating.
并不意味着不使用 Mutex,下面是一个决策树:
是否修改共享状态? │ ├── 是 → Mutex / RWMutex / atomic │ └── 否 │ └── goroutine 之间要传递数据/任务/事件? │ ├── 是 → channel └── 否 │ └── 只是等任务结束? → WaitGroup / errgroup
永远不要尝试在没有先获取该变量的互斥锁的情况下,从多个 goroutine 访问该变量。这可能导致难以追踪的奇怪错误。“15.10 用数据争用检测器查找并发问题<360页>”一节将介绍如何检测这些问题。
sync.Map——可能不是你想要的 Map
在
sync包中,有一个Map类型,它提供了一个并发安全版的映射(map)。但由于其实现中的权衡(trade-offs),sync.Map仅适用于非常特定的场景:
- 当有一个共享的 map,其中“键值对”被插入一次,但被读取多次时;
- 当 goroutine 共享该 map,但不会访问彼此的“键值”时。
此外,由于
sync.Map是在泛型(Generics)引入标准库之前添加的,它使用any作为键和值的类型。这意味着编译器无法确保是否使用了正确的数据类型。鉴于这些限制,在极少需要跨多个 goroutine 共享 map 的情况下,应使用由
RWMutex保护的内置 map。
通道实践
通道的两端分别是生产者和消费者,从通用的数据流角度来看:
- 生产者生产数据、消费者消费数据
- 由生产者负责主动停止生产数据
- 生产者停止生产数据后,消费者必须能感知到(主动感知 / 被动通知),从而让消费者停止消费,否则消费者就阻塞住了。Go 中通常通过关闭通道来实现,消费者可以根据
, ok惯用法来判断生产者是否关闭了通道。
单个生产者关闭共享通道
close(ch) 应由生产者一侧负责,而不是消费者:
func producer(ch chan<- int) {
defer close(ch)
for i := 0; i < 10; i++ {
ch <- i
}
}
多个生产者关闭共享通道
不能让每个生产者自己 close(ch),否则会重复关闭导致 panic
通常用 sync.WaitGroup,等所有 producer 都结束后,由一个协调 goroutine 关闭 channel:
ch := make(chan int)
var wg sync.WaitGroup
for i := 0; i < 3; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for j := 0; j < 5; j++ {
ch <- j
}
}(i)
}
go func() {
wg.Wait()
close(ch)
}()
for {
select {
case v, ok := <-ch:
if !ok {
fmt.Println("所有生产者都结束了")
return
}
fmt.Println(v)
}
}
消费者消费单个通道
如果只有一个 channel,其实不需要 select,直接 range 更自然:
for v := range ch {
fmt.Println(v)
}
range ch 会一直消费,直到 channel 被关闭且缓冲区被读空。
或者使用 for-select:
for {
select {
case v, ok := <-ch:
if !ok {
ch = nil
fmt.Println("生产者结束")
return
}
fmt.Println(v)
case <-ctx.Done():
return
}
}
消费者消费多个通道
消费者需要等待多个 channel 都消费完才退出时,常见写法是:某个 channel 关闭后把它设为 nil,这样后续 select 不会再选中它;直到两个都变成 nil 才退出
检测到某个通道关闭后,必须设置为 nil,否则
select可能持续选中这个已经关闭的 channel。
func consume(ch1, ch2 <-chan int) {
for ch1 != nil || ch2 != nil {
select {
case v, ok := <-ch1:
if !ok {
fmt.Println("ch1 结束")
ch1 = nil
continue
}
fmt.Println("消费 ch1:", v)
case v, ok := <-ch2:
if !ok {
fmt.Println("ch2 结束")
ch2 = nil
continue
}
fmt.Println("消费 ch2:", v)
}
}
fmt.Println("全部消费完成,退出")
}
评论