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.modgo 指令的值更改为 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 中限制操作时间时,通常会使用基于 selectcontext 的模式。这种模式是 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.OnceFuncsync.OnceValuesync.OnceValues),使得函数恰好运行一次变得更加容易。它们的唯一区别在于传入函数的返回值数量(0 个、1 个或 2 个)。sync.OnceValuesync.OnceValues 是泛型(Generics)函数,能自动适应原始函数的返回值类型(如 intstring、自定义类型等),无需手动指定类型。

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,它有两个方法:LockUnlock。当调用 Lock 时,如果此时已经有另一个 goroutine 正在执行被这个锁保护的“临界区(critical section)”(即需要互斥访问的代码块或数据),那么当前这个调用 Lock 的 goroutine 就会暂停/阻塞,直到那个在临界区的 goroutine 执行完毕并调用了 Unlock。一旦临界区空闲,锁就被当前 goroutine 获取。然后,当前 goroutine 就可以安全地执行临界区内的代码了。调用 Unlock 表示 goroutine 完成了对临界区的操作,锁被释放,其他等待的 goroutine 可以尝试获取锁,这标志着临界区的结束。
  • 第二种是 RWMutex,它允许区分读锁(独占)与写锁(共享)。任何时刻只能有一个写者进入临界区(通过 Lock 获取写锁,Unlock 释放写锁),其他写者必须等待。但是,多个读者可以同时持有读锁(通过 RLock 获取读锁,RUnlock 释放读锁)。

必须确保每次获取锁后都能正确释放锁。使用 defer 语句在调用 LockRLock 后立即调用 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.WaitGroupsync.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("全部消费完成,退出")
}

评论