并发concurrency

  • goroutine 不是简单的“线程池”。它是 Go 运行时管理的轻量级并发执行单元,由调度器把大量 goroutine 复用到少量操作系统线程上运行。它的初始栈很小,并且可以按需增长,因此创建和切换成本远低于直接创建大量系统线程。
  • 并发不是并行:Concurrency Is Not Parallelism。并发关注任务结构和调度,让多个任务能够交替推进;并行则要求多个任务在不同 CPU 核上同时执行。Go 可以通过 GOMAXPROCS 控制可同时执行 Go 代码的逻辑处理器数量。
  • Go 提倡通过通信来共享内存,而不是通过共享内存来通信。
1
2
3
4
5
6
7
8
func main() {
go Go()
// time.Sleep(2 * time.Second)
}

func Go() {
fmt.Println("Go Go Go!")
}

上面通常没有输出,因为 main goroutine 已经退出,进程随之结束。

1
2
3
4
5
6
7
8
9
func main() {
go Go()
// 这里只是演示等待 goroutine 输出,生产代码不要依赖 sleep 做同步。
time.Sleep(2 * time.Second)
}

func Go() {
fmt.Println("Go Go Go!")
}

输出:

1
Go Go Go!

更可靠的做法是使用 sync.WaitGroup 或 channel 等同步手段等待 goroutine 结束。

Channel

  • Channel 是 goroutine 沟通的桥梁,大都是阻塞同步的
  • 通过 make 创建,close 关闭
  • Channel 是引用类型
  • 可以使用 for range 来迭代不断操作 channel
  • 可以设置单向或双向通道
  • 可以设置缓存大小,在未被填满前不会发生阻塞
1
2
3
4
5
6
7
8
9
10
func main() {
// 是go一种特殊的数据类型,有点像Linux系统中的管道/消息队列
c := make(chan bool)
go func() {
fmt.Println("Go Go Go!")
c <- true
}()
// 等待从通道里读取值
<-c
}
1
2
3
4
5
6
7
8
9
10
11
12
func main() {
c := make(chan bool)
go func() {
fmt.Println("Go Go Go!")
c <- true
close(c)
}()
// 可以使用for range来迭代不断操作channel,直到关闭channel
for v := range c {
fmt.Println(v)
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
func main() {
// 有缓存的 channel,主 goroutine 发送后可能直接退出,因此不保证输出。
c := make(chan bool, 1)
go func() {
fmt.Println("Go Go Go!")
<-c
}()
c <- true
}

func main() {
// 无缓存的 channel 会在发送/接收两端同步,输出:Go Go Go!
c := make(chan bool)
go func() {
fmt.Println("Go Go Go!")
<-c
}()
c <- true
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
// 在并发情况下,程序不是按 goroutine 创建顺序执行。
// index=9 执行完不代表所有 goroutine 都执行完。
func main() {
runtime.GOMAXPROCS(runtime.NumCPU())
c := make(chan bool)
for i := 0; i < 10; i++ {
go Go(c, i)
}
<-c
}

func Go(c chan bool, index int) {
a := 1
for i := 0; i < 100000000; i++ {
a += i
}
fmt.Println(index, a)

if index == 9 {
c <- true
}
}

输出(每次都会不同):
0 4999999950000001
3 4999999950000001
1 4999999950000001
2 4999999950000001
6 4999999950000001
7 4999999950000001
9 4999999950000001
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
// 解决上面不能判断所有 goroutine 是否执行完的问题:
// 通过创建带缓冲区的 channel 收集完成信号。
func main() {
runtime.GOMAXPROCS(runtime.NumCPU())
c := make(chan bool, 10)
for i := 0; i < 10; i++ {
go Go(c, i)
}
for i := 0; i < 10; i++ {
<-c
}
}

func Go(c chan bool, index int) {
a := 1
for i := 0; i < 100000000; i++ {
a += i
}
fmt.Println(index, a)

c <- true
}

输出(每次都能输出10个):
4 4999999950000001
0 4999999950000001
6 4999999950000001
9 4999999950000001
1 4999999950000001
7 4999999950000001
8 4999999950000001
5 4999999950000001
2 4999999950000001
3 4999999950000001
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 多个 goroutine 协作通过同步包来实现。
func main() {
runtime.GOMAXPROCS(runtime.NumCPU())
wg := sync.WaitGroup{}
wg.Add(10)
for i := 0; i < 10; i++ {
go Go(&wg, i)
}
wg.Wait()
}

func Go(wg *sync.WaitGroup, index int) {
a := 1
for i := 0; i < 100000000; i++ {
a += i
}
fmt.Println(index, a)

wg.Done()
}

Select

  • 可处理一个或多个 channel 的发送与接收
  • 同时有多个可用的 case 时,select 会伪随机选择一个执行
  • 可用空的 select 来阻塞 main 函数
  • 可设置超时
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
func main() {
c1, c2 := make(chan int), make(chan string)
go func() {
for {
select {
case v, ok := <-c1:
if !ok {
break
}
fmt.Println("c1", v)
case v, ok := <-c2:
if !ok {
break
}
fmt.Println("c2", v)
}
}
}()

c1 <- 1
c2 <- "hi"
c1 <- 3
c2 <- "hello"

close(c1)
close(c2)
}

输出:
c1 1
c2 hi
c1 3
c2 hello
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
// c1、c2有一个关闭程序就会退出
// 两个都关闭才退出是不可行的,我们只有对其中一个进行判断
func main() {
c1, c2 := make(chan int), make(chan string)
o := make(chan bool)
go func() {
for {
select {
case v, ok := <-c1:
if !ok {
o <- true
break
}
fmt.Println("c1", v)
case v, ok := <-c2:
if !ok {
o <- true
break
}
fmt.Println("c2", v)
}
}
}()

c1 <- 1
c2 <- "hi"
c1 <- 3
c2 <- "hello"

close(c1)
close(c2)

<-o
}

输出:
c1 1
c2 hi
c1 3
c2 hello
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
// 通过select来作为一个发送者的应用
func main() {
c := make(chan int)
go func() {
for v := range c {
fmt.Println(v)
}
}()

for {
select {
case c <- 0:
case c <- 1:
}
}
}

输出:
0
1
1
1
0
1
1
1
0
0
0
1
.
.
.
1
2
3
4
5
6
7
8
9
10
11
12
13
// 可设置Select超时
func main() {
c := make(chan bool)
select {
case v := <-c:
fmt.Println(v)
case <-time.After(3 * time.Second):
fmt.Println("Timeout")
}
}

输出:
Timeout

思考问题

  • 创建一个 goroutine,与 main goroutine 按顺序相互发送信息若干次并打印
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
var c chan string

func Pingpong() {
i := 0
for {
fmt.Println(<-c)
c <- fmt.Sprintf("From Pingpong: Hi, #%d", i)
i++
}
}

func main() {
c = make(chan string)
go Pingpong()
for i := 0; i < 10; i++ {
c <- fmt.Sprintf("From main: Hello, #%d", i)
fmt.Println(<-c)
}
}

输出:
From main: Hello, #0
From Pingpong: Hi, #0
From main: Hello, #1
From Pingpong: Hi, #1
From main: Hello, #2
From Pingpong: Hi, #2
From main: Hello, #3
From Pingpong: Hi, #3
From main: Hello, #4
From Pingpong: Hi, #4
From main: Hello, #5
From Pingpong: Hi, #5
From main: Hello, #6
From Pingpong: Hi, #6
From main: Hello, #7
From Pingpong: Hi, #7
From main: Hello, #8
From Pingpong: Hi, #8
From main: Hello, #9
From Pingpong: Hi, #9