Blog

Goroutine 与 Go 调度器

goroutine 是和程序其余部分同时运行的函数,Go 可以同时跑成千上万个。学会启动它、等待它、避免泄漏,并看懂调度器怎样让少数线程承载它们。

goroutine 是和程序其余部分同时运行的函数。在函数调用前面写上 go,就启动了一个。这一步一分钟就能学会。难的是弄清 goroutine 什么时候结束,以及在这期间运行时拿它们做了什么。

本文讲怎样启动 goroutine、怎样等待它们、它们为什么便宜,以及运行它们的调度器。最后讲 goroutine 最常见的 bug:泄漏。下面每个程序都在 Go 1.26 上跑过,输出直接从运行结果粘贴而来。

go 启动 goroutine,main 不会等它

在函数调用前加上 go,这个调用就会在一个新的 goroutine 里开始执行,而当前代码立刻往下走,不等调用结束。

package main

import (
	"fmt"
	"time"
)

func report() {
	time.Sleep(time.Minute) // stands in for slow work
	fmt.Println("report finished")
}

func main() {
	go report()
	fmt.Println("main done")
}

输出:

main done

报告一直没有打印出来。main 启动了 goroutine,打印自己那一行,然后就返回了,离一分钟还早得很。main 一返回,程序就结束。Go 不会等其他 goroutine,也不会提醒你还有 goroutine 在跑。它们直接被停掉。

这里放 sleep,是为了让每次运行的结果都一样。去掉它,改写成 go fmt.Println("hello from the goroutine"),结果就变成了抛硬币。我们在一台机器上把这个版本跑了 200 次,大约一半打印出了 goroutine 那一行,另一半没有。程序里没有任何东西决定是哪一种。一个只在 goroutine 碰巧够快时才正常的程序,就是有 bug。

你有时会看到有人在 main 末尾加一句 time.Sleep 来“修”这个问题。在空闲的机器上它能把问题藏起来,机器一忙就失灵。真正的做法是等 goroutine 自己报告已经完成。

sync.WaitGroup 正确地等待

sync.WaitGroup 是一个计数器,记着还有多少 goroutine 在运行;Wait 会一直阻塞,直到计数归零。

package main

import (
	"fmt"
	"sync"
)

func main() {
	names := []string{"ana", "bo", "chen"}
	greetings := make([]string, len(names))

	var wg sync.WaitGroup
	for i, name := range names {
		wg.Add(1)
		go func() {
			defer wg.Done()
			greetings[i] = "hello, " + name
		}()
	}
	wg.Wait()

	for _, g := range greetings {
		fmt.Println(g)
	}
}

输出:

hello, ana
hello, bo
hello, chen

要掌握三个调用:

  • wg.Add(1) 把计数器加一。要在 go 语句之前调用,而且是在负责等待的那个 goroutine 里调用。如果放到新 goroutine 里面,Wait 可能先执行,看到计数为零,就提前返回了。
  • wg.Done() 把计数器减一。把它放在 defer 后面,函数提前返回时它也会执行。
  • wg.Wait() 一直阻塞,直到计数器回到零。

注意输出一直是按顺序的。三个 goroutine 的运行顺序不固定,但每个都只写自己的格子 greetings[i]。没有两个 goroutine 碰同一个元素,而 main 要等 Wait 返回后才读这个切片。打印集中在一个地方,按固定顺序进行。这个习惯能让并发程序的输出保持可预测。

Go 1.25 新增了 wg.Go

从 Go 1.25 开始,WaitGroup 有了 Go 方法,替你完成 Addgo 语句和 Done

package main

import (
	"fmt"
	"sync"
)

func main() {
	squares := make([]int, 5)

	var wg sync.WaitGroup
	for i := range squares {
		wg.Go(func() {
			squares[i] = i * i
		})
	}
	wg.Wait()

	fmt.Println(squares)
}

输出:

[0 1 4 9 16]

wg.Go(f) 等同于先 wg.Add(1),再启动一个 goroutine,调用 f 之后执行 Done。你不会忘了 Done,也不会把 Add 放错位置。传进去的函数不接收参数、也没有返回值,所以它要通过闭包拿到需要的值,就像这里的 i。本文后面都用 wg.Go。不过现有代码里大多还是 AddDone,两种写法都要能看懂。

goroutine 里的闭包与循环变量

在循环里启动的 goroutine 几乎总会用到循环变量。讲控制流的那一部分解释过,这里过去是个陷阱。

看上面的 squares[i] = i * i。循环走到这一行时,goroutine 并没有马上运行。它会稍晚一点才运行,可能循环都已经结束了。那它看到的是哪个 i

从 Go 1.22 开始,每次迭代都有自己的 i。第三轮启动的 goroutine 看到的是第三轮的 i,之后也没有任何东西会改它。所以这个程序是对的。

Go 1.22 之前,整个循环共用一个 i。运行得晚的 goroutine 读到的都是循环当时走到的值,往往是最后一个。为老版本 Go 写的代码会绕开这个问题:在循环里写 i := i,或者把值作为参数传进去,go func(i int) { ... }(i)。在 Go 1.22 及以后,两种都不需要。这条规则看的是 go.mod 里的 go 那一行,而不是你安装的 Go 版本,所以一个仍然写着 go 1.21 的模块还是老行为。

goroutine 很便宜

goroutine 的开销远小于操作系统线程,所以启动几千个是 Go 里的常规操作。

package main

import (
	"fmt"
	"sync"
)

func main() {
	const n = 10_000
	results := make(chan int, n)

	var wg sync.WaitGroup
	for i := range n {
		wg.Go(func() {
			results <- i + 1
		})
	}
	wg.Wait()
	close(results)

	total := 0
	for r := range results {
		total += r
	}
	fmt.Println("goroutines:", n)
	fmt.Println("total:", total)
}

输出:

goroutines: 10000
total: 50005000

一万个 goroutine 各往通道(channel)里发送一个数。通道能装下全部一万个值,所以没有哪个发送方需要等待。Wait 之后,main 关闭通道,把里面的值全部加起来。总和是 1 + 2 + … + 10,000,也就是 50,005,000,不管 goroutine 按什么顺序运行都一样。

为什么这么便宜?讲内存的那一部分演示过,新 goroutine 的栈一开始只有几 KB,需要时才增长。操作系统线程通常一开始就预留一块大得多的固定栈。而且创建 goroutine 是 Go 运行时自己的事,不用去找操作系统。

便宜不等于免费。每个 goroutine 在结束之前,都占着自己的栈和它引用的所有东西。一万个很快结束的 goroutine 没问题。一万个永远不结束的 goroutine 就是泄漏,本文最后会讲。

这里的通道承担的工作,到讲通道的那一部分会正式介绍。现在先把 results <- i + 1 理解成“把这个值放进队列”,把 for r := range results 理解成“一直取值,直到队列关闭”。

并发不是并行

并发是指程序被组织成几个各自推进的任务。并行是指几个任务在同一时刻执行,分别在不同的 CPU 核心上。goroutine 给你的是并发。能不能同时得到并行,取决于运行时被允许使用多少个核心。

这个上限叫 GOMAXPROCS,也就是同一时刻能运行 Go 代码的操作系统线程数。runtime.NumCPU() 报告进程能使用多少个逻辑 CPU,GOMAXPROCS 默认就从这个数开始。从 Go 1.25 开始,在 Linux 上,默认值还会遵守 cgroup 的 CPU 限制,也就是容器设置的那种。一个被限制在两个 CPU 的程序,得到的 GOMAXPROCS 会比机器的核心数小。这两个数都取决于程序在哪里运行,所以这里的程序都不打印它们。

你可以自己设置 GOMAXPROCS。设成 1 就完全没有并行了,而 goroutine 照样能工作:

package main

import (
	"fmt"
	"runtime"
	"sync"
)

func main() {
	runtime.GOMAXPROCS(1)

	var mu sync.Mutex
	total := 0

	var wg sync.WaitGroup
	for i := range 1000 {
		wg.Go(func() {
			mu.Lock()
			total += i
			mu.Unlock()
		})
	}
	wg.Wait()

	fmt.Println("GOMAXPROCS:", runtime.GOMAXPROCS(0))
	fmt.Println("total:", total)
}

输出:

GOMAXPROCS: 1
total: 499500

runtime.GOMAXPROCS(1) 设置上限,runtime.GOMAXPROCS(0) 只读取、不修改。上限为 1 时,任一时刻只有一个 goroutine 在运行 Go 代码。一千个 goroutine 轮流执行,结果仍然是 0 + 1 + … + 999。

互斥锁还是需要的。轮流执行并不代表每个 goroutine 都会在下一个开始前做完加法,而且程序本来就不该依赖这个设置。讲 sync 的那一部分会介绍互斥锁和竞态检测器。实际代码里你很少需要设置 GOMAXPROCS,默认值通常就是对的。

调度器怎样运行 goroutine

Go 运行时在少数几个操作系统线程上运行大量 goroutine,决定哪个 goroutine 在哪里运行的,就是调度器。

用十岁孩子能懂的话说

想象一个餐厅厨房,厨师很多,灶台只有几个。厨师是 goroutine,灶台是真正干活的线程。可能有五十个厨师,只有四个灶台。

主厨就是调度器,他决定哪个厨师站到哪个灶台前。每个厨师做一会儿菜,然后换另一个厨师上。

有时厨师得等烤箱。他不会站在灶台前干等,而是让开位置,主厨再派别人到那个灶台。烤箱“叮”一声响了,等待的厨师重新排队等灶台,排到的甚至可能不是他原来那个。

所以四个灶台能让五十个厨师都有活干。

准确的说法

运行时的调度器用到三种对象,通常叫 GMP

  • G 是 goroutine:它的栈,以及它执行到了哪里。
  • M 是 machine,也就是操作系统线程。只有 M 能真正执行代码。
  • P 是 processor:运行 Go 代码的许可,外加一个本地队列,里面是准备好运行的 goroutine。P 的数量正好是 GOMAXPROCS 个。

M 必须持有一个 P 才能运行 Go 代码。M 从这个 P 的运行队列里取一个 G 来运行。这叫 M:N 调度:许多 goroutine 共用较少的操作系统线程,由 Go 运行时而不是操作系统决定哪个 goroutine 在哪个线程上运行。

M + P 运行中 本地运行队列 M0 P0 M1 P1 等待中 G1 挂起在 <-ch G2 G3 G4 G5 G1 G1 在 P0(线程 M0)上运行,G2 在 P1(线程 M1)上。其余在队列里等。 G1 执行 <-ch,没有值可收。它被挂起,P0 空出来了。 M0 不闲着:P0 从运行队列取出 G3 来运行。 G2 往 ch 发送。G1 重新可运行,进入 P1 的运行队列。 G2 结束。P1 运行 G1,这次在线程 M1 上,而不是 M0。

简化示意:两个 P,各自绑定一个线程,各自有一个本地运行队列。G1 在通道上阻塞时被挂起,交出它的 P,于是线程去运行下一个 goroutine,而不是干等。G1 变为可运行后回到某个运行队列,这个队列可能属于另一个 P。真实的调度器还有一个全局运行队列,空闲的 P 也会从忙碌的 P 那里窃取任务。

下面用文字把这几步再说一遍,以防动画无法播放:

  1. P0 绑定线程 M0,运行 G1。P1 绑定线程 M1,运行 G2。G3 和 G4 在 P0 的队列里等待,G5 在 P1 的队列里。
  2. G1 试图从一个空通道接收。调度器把 G1 挂起,进入等待状态。它不再占用线程,也不再占用 P。
  3. M0 不等 G1。P0 从队列里取出下一个 goroutine G3 来运行。
  4. G2 往那个通道发送了一个值。G1 因此变为可运行,进入 P1 的队列,因为发送发生在 P1 上。
  5. G2 结束,P1 运行 G1。G1 最初在线程 M0 上运行,现在跑在 M1 上。

goroutine 不绑定某个线程。哪个 P 接手它,它就在哪里运行。

调度器还要处理三件事:

  • 阻塞的系统调用。有些工作,比如读文件,会在操作系统内部阻塞整个线程。发生这种情况时,运行时会把 P 从被阻塞的 M 那里拿走,交给另一个线程,让其他 goroutine 继续运行。runtime 包的文档说得很直白:阻塞在系统调用里的线程不计入 GOMAXPROCS 的限制。网络读取不一样。运行时用网络轮询器(network poller)等待 socket,所以等待网络的 goroutine 会像 G1 一样被挂起,不占用线程。
  • 抢占。goroutine 不能一直霸占一个 P。从 Go 1.14 开始,运行时在大多数平台上都能中断 goroutine,哪怕它在一个没有函数调用的紧凑循环里,所以一个忙碌的 goroutine 饿不死其他 goroutine。
  • 任务窃取。P 的队列空了,就先去全局运行队列里找,再从其他 P 那里拿 goroutine。为了看得清楚,动画省略了这一步。

这个比喻的局限:厨房里站在灶台前的只有一种角色,就是厨师。Go 在这里有两个不同的概念:线程(M),以及运行 Go 代码的许可(P)。卡在系统调用里的线程,就像一个灶台前的厨师动弹不得,运行时的处理办法是再搬来一个灶台,也就是一个新线程,把 P 交给它。另外,主厨每做一个决定都会动脑子。Go 调度器只是飞快地做同样简单的选择,它也不知道哪个 goroutine 更重要。

goroutine 泄漏

goroutine 泄漏指的是一个永远不会结束的 goroutine,通常是因为它在等一个再也没人会用的通道。

package main

import (
	"fmt"
	"runtime"
)

func firstResult() int {
	ch := make(chan int)
	go func() {
		ch <- 42 // nobody will ever receive this
	}()
	return -1 // gave up without reading ch
}

func main() {
	fmt.Println("before:", runtime.NumGoroutine())
	for range 3 {
		firstResult()
	}
	fmt.Println("after:", runtime.NumGoroutine())
}

输出:

before: 1
after: 4

runtime.NumGoroutine() 报告当前有多少个 goroutine。循环之前只有一个,就是 main 本身。每次调用 firstResult 都会启动一个 goroutine,它试图往一个无缓冲通道发送。在无缓冲通道上发送,要等到有人接收才会完成。可是 firstResult 没有接收就返回了,其他代码也拿不到 ch。于是每个 goroutine 都永远等下去。调用三次,卡住三个 goroutine,计数就是 4。

程序不会崩溃。只有所有 goroutine 都卡住时,Go 才会报告死锁,而这里 main 还在继续运行。垃圾回收器也帮不上忙。被阻塞的 goroutine 它不会释放,哪怕已经没有任何东西能访问到那个通道。如果一个服务器每个请求调用一次 firstResult,这个计数就会随请求不断上涨,直到内存耗尽。

这里的计数是精确的,因为 goroutine 从 go 创建它的那一刻就被计入,而这三个永远不会退出。在实际代码里,发现泄漏的办法也是数 goroutine。持续观察 runtime.NumGoroutine(),或者查看 net/http/pprof 提供的 goroutine profile。如果负载没变,数字却一直涨,就说明有东西在泄漏。一长串 goroutine 都卡在同一行,就告诉了你泄漏在哪里。

Go 1.26 还带了一个实验性的 profile,叫 goroutineleak。它只存在于用 GOEXPERIMENT=goroutineleakprofile 构建的程序里,借助垃圾回收器找出那些阻塞在任何运行中 goroutine 都访问不到的东西上的 goroutine。我们在这个程序上试了一下,结果很好地提醒了“并发”意味着什么:有一次运行它报告的是两个泄漏的 goroutine,而不是三个。最可能的原因是,取 profile 时第三个已经创建,但还没执行到发送那一步。泄漏检测器只能看到已经卡住的 goroutine。

这个 bug 通常有两种修法:

  • 给通道留出放一个值的空间,make(chan int, 1)。发送立刻完成,即使没人读,goroutine 也会退出。
  • 给 goroutine 一个听到“停下”的途径。这正是 selectcontext 的用途,分别在讲通道和讲 context 的部分介绍。

要记住的规则:启动 goroutine 之前,先想清楚它会怎样结束。

接下来

本文用了通道和互斥锁,但没有解释它们。讲通道和 select 的那一部分介绍 goroutine 之间怎样传值,包括无缓冲通道、关闭通道和 select。讲 sync 的那一部分介绍互斥锁、竞态检测器和 atomic 包,只要 goroutine 共享内存,你就会用到它们。

要点

  • go f() 在新的 goroutine 里启动 f,不会等待。main 返回时程序结束,还在运行的 goroutine 会被停掉。
  • sync.WaitGroup 等待。从 Go 1.25 开始,wg.Go(func() { ... }) 替你完成 AddDone
  • goroutine 很便宜,几千个很正常。通过通道或受保护的切片收集结果,并集中在一个地方打印。
  • 并发说的是程序的结构。并行说的是同一时刻真正在执行,GOMAXPROCS 限制它。
  • 调度器通过处理器(P)把 goroutine(G)放到线程(M)上运行。阻塞的 goroutine 会被挂起,它的线程转去做别的工作。
  • 从 Go 1.22 开始,每次循环迭代都有自己的变量,所以在循环里启动的 goroutine 看到的是本轮迭代的值。
  • 永远阻塞在通道上的 goroutine 就是泄漏。你启动的每个 goroutine,都要知道它会怎样结束。

启动 goroutine 只要一个词。弄清它什么时候结束,是你的责任。

这篇文章对你有帮助吗?

点一颗爱心来评分!

平均评分 0 / 5. 投票总数: 0

还没有人投票。来做第一个评分的人吧。