12. 并发

Go 的并发(Concurrency)是 Go 最具特色、也是最重要的特性之一。并不像其他编程语言,如 Java,C++ 等中普遍采用的多进程、多线程方式,也不像 Nodejs 中采用的基于回调的异步 I/O 方式,在 Go 语言中提供了基于协程(Coroutine),协程是一种更轻量化的线程。

12.1. 并发与并行

首先,我们需要弄清楚并发于并行的区别,这便于我们更好理解 Go 语言在并发编程中的独特性。简单来讲,并发是是指系统有能力在重叠的时间段内处理多个任务,但并不一定在同一个时刻执行,通常是指在单核的场景下。而并行是指系统在同一个物理时刻真正同时执行多个任务,通常是指在多核的场景下。两者的区别可以下图表示:

在众多的并发方案中,有以下几种主流的技术方案,如多进程,多线程,基于回调的异步 I/O以及协程。简单来说,协程就是更轻量化的线程,优点是编程简单,缺点是需要语言的支持。Go 语言中就选择使用协程的方案来实现并发。

12.2. Goroutine 和 sync

在 Go 语言中,内置了支持并发的能力——协程,在 Go 语言中称为 goroutine。

12.2.1. Goroutine

goroutine 是轻量级线程,goroutine 的调度是由 Golang 运行时进行管理的。goroutine 语法格式如下所示:

go 函数名(参数列表)

我们以一个简单的例子开始,文件命名为 main.go,其中的代码如下:

package main

import "fmt"

func Say(i int) {
	fmt.Printf("coroutine: %d\n", i)
}

func main() {
	for i:=0; i<10; i++ {
		go Say(i)
	}
}

上述代码中,通过使用 go 语句开启一个 goroutine。但是在执行 go run main.go 时,发现本应该打印的文字并没有显示。我们发现 Say() 函数还未执行,主程序已经退出了,因此压根就不会输出任何个文字。这时候需要主程序在等待所有 goroutine 执行完成后,主程序再退出,在 Go 语言中,提供了 sync 和 channel 两种方式支持 goroutine 的并发。

12.2.2. sync

sync 是 Go 标准库中专门用于并发同步的包,当有多个 goroutine 同时运行时,sync 提供了一组工具,协调这些 goroutine 的执行顺序。WaitGroup 正是用来解决多个 goroutine 之间的同步问题的,WaitGroup 可以理解为 goroutine 的任务计数器,为此,提供了三个重要的方法,分别为:

  • Add():作用是修改 WaitGroup 内部的任务计数器
  • Done():作用是当前有一个 goroutine 完成了任务
  • Wait():作用是阻塞当前的 goroutine,等待 WaitGroup 计数器变为 0

针对上面的例子,增加 WaitGroup 后的代码如下:

package main

import (
	"fmt"
	"sync"
)

// 声明 WaitGroup
var wg sync.WaitGroup

func Say(i int) {
	fmt.Printf("coroutine: %d\n", i)
	// goroutine 完成
	wg.Done()
}

func main() {
	for i:=0; i<10; i++ {
		// 修改计数器
		wg.Add(1)
		go Say(i)
	}
    // 等待计数器为 0
	wg.Wait()
}

此时,再运行该程序,我们发现每个 goroutine 可以正常输出了。但是上述的代码还有个细微的问题,一般在 Say() 函数中,我们通常采用如下的方式:

func Say(i int) {
	// goroutine 完成
	defer wg.Done()
	fmt.Printf("coroutine: %d\n", i)
}

这里 defer 是 Go 语言中一个重要的关键字,其作用是把一个函数调用推迟到当前函数即将返回时执行。用在这里的目的是防止在 Say() 函数中异常退出,导致该 goroutine 一直无法结束,最终导致程序报错,使用 defer 保证在该 goroutine 退出时标记该 goroutine 结束。

在 sync 包中除了 WaitGroup 外,还提供了对锁的支持,如 MutexRWMutex等,这部分内容将在后面介绍到。

12.3. channel

12.3.1. channel 的基本使用

channel 是 Go 语言中另一种支持并发的方式,实际上,channel 是 Go Runtime 提供的一种并发安全的数据传递机制,它允许多个 goroutine 通过发送和接收数据进行同步和通信。从这个定义可以看出 channel 是支持 goroutine 发送和接收数据的容器。

在使用 channel 进行多个 goroutine 通信前,需要先创建 channel,可以使用 make 函数创建,创建的语法如下:

make(chan T) // 无缓冲的 channel
make(chan T, capacity) // 有缓冲的 channel,容量大小为 capacity

其中,T 表示 channel 中传递的数据类型。从创建的语法来看,channel 是类型安全的,在创建的时候必须指定传递的数据的理性。有了 channel 后,还有两个操作,一个是向 channel 发送数据,另一个是从 channel 取出数据,使用方法如下代码所示:

// 创建 chan int 类型的 channel
var ch chan int = make(chan int)
// 向 channel 发送数据
ch <- 1
// 从 channel 取出数据
msg := <-ch

上述例子的代码用 channel 也能实现,如下代码:

package main

import (
	"fmt"
)

// 创建 channel
var ch chan int = make(chan int)

func Say(i int) {
	fmt.Printf("coroutine: %d\n", i)
	// 向 channel 发送数据
	ch <- i
}

func main() {
	for i:=0; i<10; i++ {
		go Say(i)
	}

	for i:=0; i<10; i++ {
		// 从 channel 接收数据
		msg := <-ch
		fmt.Println(msg)
	}
}

上面的代码所表示的过程可以由下图表示:

从上面也可以看出,简单来说,channel 就是队列,不同的是,channel 的队列是并发安全的队列。实际上 channel 的功能还远不止这些。

12.3.2. channel 是有方向的

从上面的代码可以看到,我们可以向 channel 发送数据,也可以从 channel 接收数据,代码如下:

// 向 channel 发送数据
ch <- i
// 从 channel 接收数据
msg := <-ch

默认的 channel 是双向的,同样,我们可以限制 channel 的方向,如下:

// 只能发送
make(chan<- int)
// 只能接收
make(<-chan int)

有无缓冲

在上面的创建语法来看,channel 的创建分为有没有缓冲,代码如下:

make(chan T) // 无缓冲的 channel
make(chan T, capacity) // 有缓冲的 channel,容量大小为 capacity

其中,在无缓冲模式下,发送和接收必须同时发生,如果只有发送,此时并不会向 channel 中写入数据,发送的操作会被阻塞,代码如下:

// 创建 channel
var ch chan int = make(chan int)

func Say(i int) {
	// 向 channel 发送数据
	ch <- i
	fmt.Printf("coroutine: %d\n", i)
}

func main() {
	// 只有发送
	for i:=0; i<10; i++ {
		go Say(i)
	}

	time.Sleep(5 * time.Second)
}

运行代码,我们发现程序并不会向 channel 中发送数据。而对于有缓冲的 channel,在缓冲未满的情况下会一直向 channel 中写入数据,代码如下:

// 创建带有容量为 5 的 channel
var ch chan int = make(chan int, 5)

func Say(i int) {
	// 向 channel 发送数据
	ch <- i
	fmt.Printf("coroutine: %d\n", i)
}

func main() {
	// 只有发送
	for i:=0; i<10; i++ {
		go Say(i)
	}

	time.Sleep(5 * time.Second)
}

我们发现,在填满容量大小后,后面的写入就阻塞了。

12.3.3. channel 中的数据满足先进先出

在 channel 中,数据严格按照先进先出的顺序接收,以下面的代码为例:

// 创建带有容量为 5 的 channel
var ch chan int = make(chan int, 5)

func Say(i int) {
	// 向 channel 发送数据
	fmt.Printf("coroutine: %d\n", i)
	ch <- i
}

func main() {
	// 只有发送
	for i:=0; i<10; i++ {
		go Say(i)
	}

	for i:=0; i<10; i++ {
		msg := <-ch
		fmt.Printf("read: %d\n", msg)
	}
}

运行后,得到的结果为:

coroutine: 0
read: 0
coroutine: 4
read: 4
coroutine: 1
read: 1
coroutine: 2
read: 2
coroutine: 8
read: 8
coroutine: 9
coroutine: 3
coroutine: 6
coroutine: 7
read: 9
read: 3
read: 6
read: 7
coroutine: 5
read: 5

接收数据的顺序与写入的顺序一致。

12.3.4. channel 是并发安全的

这个也是最重要的特性,无论是发送数据还是接收数据,channel 都是并发安全的。以下面的代码为例:

func Worker(i int, ch chan int) {
	for msg := range ch {  
		fmt.Printf("worker: %d, ch message: %d\n", i, msg)
	}
}

func main() {
	var ch chan int = make(chan int, 20)
	for i:=0; i<10; i++ {
		ch <- i
	}
	close(ch)

	for i:=0; i<3; i++ {
		go Worker(i, ch)
	}

	// 不让主进程退出
	time.Sleep(5 * time.Second)
}

channel 中数据每次只会消费一次。

12.3.5. range 和 close

上面的代码中也使用到了 for ... range 的语法从 channel 中接收数据,for ... range 的功能是不断从 channel 中接收数据,语法相当于下面的代码:

for {
    msg, ok := <-ch
    if !ok {   // 如果 channel 已关闭且为空,退出循环
        break
    }
    // 处理 msg
}

此时,就需要我们在发送端使用 close() 函数主动关闭 channel。

12.4. Mutex

在 sync 中除了提供 WaitGroup 来解决 goroutine 之间的同步问题之外,还提供锁的相关操作,包括 Mutex,主要用于解决多个 goroutine 如何安全访问共享数据。

Mutex 的核心思想非常简单,同一时刻,只允许一个 goroutine 进入临界区。Mutex 的使用方法也比较简单,如下:

var mu sync.Mutex

mu.Lock()

// 临界区
// 访问共享数据

mu.Unlock()

具体的使用方法如下代码所示:

// 共享数据
var count int
var mu sync.Mutex

func main() {
	var wg sync.WaitGroup

	for i := 0; i < 1000; i++ {
		wg.Add(1)

		go func() {
			defer wg.Done()

			mu.Lock()
			// 访问共享数据
			count++
			mu.Unlock()
		}()
	}
	wg.Wait()

	fmt.Println(count)
}

12.5. 本章小结

在并发编程中,Go 语言并没有沿着线程的路线,而是在语言层面提供了对协程的支持,同时,Go 语言中的协程也非常简单,同时,为了解决协程同步与通信,提供了 sync 和 channel 两个工具,非常方便。