Web

Goroutine、Channel 与并发控制

Go 通过 Goroutine 执行并发任务,并使用 Channel、Context 和 sync 包中的工具协调这些任务。本文从启动一个 Goroutine 开始,逐步介绍任务等待、数据传递、取消和共享状态保护。

启动 Goroutine

普通函数调用会按照当前控制流同步执行。在调用前添加 go,可以让函数在新的 Goroutine 中执行:

go
func printMessage() {
	fmt.Println("hello")
}

func main() {
	go printMessage()
}

并发表示多个任务可以交错推进,并不保证它们在同一时刻并行运行。是否真正并行还取决于 CPU 核心数和 Go 运行时的调度。

上面的程序也不能保证输出 hello。如果 main() 先结束,整个进程会直接退出,不会继续等待其他 Goroutine。启动并发任务时,还需要明确由谁等待和结束它。

使用 WaitGroup 等待任务

只需要等待一组任务全部完成时,可以使用 sync.WaitGroup

go
var group sync.WaitGroup

group.Add(2)

go func() {
	defer group.Done()
	fmt.Println("task 1")
}()

go func() {
	defer group.Done()
	fmt.Println("task 2")
}()

group.Wait()

Add(2) 表示需要等待两个任务,每次调用 Done() 都会把计数减一,Wait() 则阻塞到计数归零。

Add() 应在启动 Goroutine 之前调用。Done() 通常放在 defer 中,确保函数从正常返回路径退出时仍会更新计数。两个任务的完成顺序并不确定,如果程序依赖固定顺序,就不应仅仅把它们放进不同的 Goroutine。

使用 Channel 传递数据

Channel 用于在 Goroutine 之间传递指定类型的数据。常见写法如下:

写法作用
chan T表示传递 T 类型数据的 Channel
make(chan T)创建无缓冲 Channel
make(chan T, n)创建容量为 n 的有缓冲 Channel
ch <- value向 Channel 发送值
value := <-ch从 Channel 接收值
close(ch)表示发送方不会再发送新值
for value := range ch持续接收,直到 Channel 关闭并被读空

无缓冲 Channel 不保存数据。发送操作会等待接收方就绪,接收操作也会等待发送方就绪,因此它既传递数据,也同步两个 Goroutine:

go
messages := make(chan string)

go func() {
	messages <- "done"
}()

message := <-messages
fmt.Println(message)

有缓冲 Channel 可以暂存一定数量的值。缓冲区写满时,新的发送操作会等待;缓冲区为空时,接收操作会等待:

go
messages := make(chan string, 2)

messages <- "first"
messages <- "second"

fmt.Println(<-messages) // first
fmt.Println(<-messages) // second

Channel 按值成功发送的先后顺序进行接收,可以理解为先进先出(FIFO)。但多个 Goroutine 的执行顺序并不确定,因此并发发送时,哪个值先进入 Channel 仍取决于实际调度。

函数参数还可以限制 Channel 的操作方向:

go
func produce(output chan<- string) {
	output <- "done"
}

func consume(input <-chan string) {
	fmt.Println(<-input)
}

chan<- string 只能发送,<-chan string 只能接收。方向限制可以让函数的职责更清楚。

使用 Channel 构建任务队列

多个 Goroutine 可以从同一个 Channel 接收任务。下面的程序启动两个 Worker,处理五个任务:

go
jobs := make(chan int, 3)

var group sync.WaitGroup

for workerID := 1; workerID <= 2; workerID++ {
	group.Add(1)

	go func(id int) {
		defer group.Done()

		for jobID := range jobs {
			fmt.Printf("worker %d handles job %d\n", id, jobID)
		}
	}(workerID)
}

for jobID := 1; jobID <= 5; jobID++ {
	jobs <- jobID
}

close(jobs)
group.Wait()

可以按以下步骤理解这段代码:

  1. jobs := make(chan int, 3) 创建容量为 3 的任务队列。缓冲区未满时,发送任务不会阻塞;缓冲区满后,发送方会等待 Worker 取走任务。
  2. for workerID := 1; workerID <= 2; workerID++ 启动两个 Worker。每个 Worker 用 for jobID := range jobs 从同一个 Channel 接收任务,直到 jobs 关闭并被读空。
  3. group.Add(1) 在启动 Worker 前增加等待计数,defer group.Done() 保证 Worker 退出时减少计数,group.Wait() 等待两个 Worker 都结束。
  4. 发送完 5 个任务后,发送方调用 close(jobs)。Worker 会先读空剩余任务,再结束 range 循环;如果不关闭 Channel,range 会一直等待新任务。

Channel 的关闭

只有发送方确定不会再写入数据,并且接收方需要感知结束时,才需要关闭 Channel。重复关闭 Channel 或向已关闭的 Channel 发送数据都会触发 panic。

使用 select 等待多个操作

select 可以同时等待多个 Channel 操作,并执行最先就绪的分支,并略过其余的分支:

go
select {
case message := <-messages:
	fmt.Println(message)
case <-time.After(time.Second):
	fmt.Println("timeout")
}

如果一秒内没有收到消息,time.After() 返回的 Channel 会就绪,程序进入超时分支。

使用 Context 传递取消信号

仅让等待方超时还不够,正在执行的 Goroutine 也应有机会结束。context.Context 可以在调用链中传递取消信号和截止时间:

go
ctx, cancel := context.WithTimeout(
	context.Background(),
	500*time.Millisecond,
)
defer cancel()

result := make(chan string, 1)

go func() {
	select {
	case <-time.After(time.Second):
		result <- "completed"
	case <-ctx.Done():
		return
	}
}()

select {
case value := <-result:
	fmt.Println(value)
case <-ctx.Done():
	fmt.Println("timed out")
}

Context 在 500 毫秒后超时,外层 select 不再等待结果,内部 Goroutine 也会收到同一个取消信号并退出。创建可取消的 Context 后,应在不再需要它时调用 cancel()

函数或方法作用
context.Background()返回一个不会被取消的空 Context,通常作为根节点使用
context.TODO()暂时不确定使用哪种 Context 时占位,正式代码中应尽快替换
context.WithCancel(parent)基于父 Context 创建可取消的 Context,并返回 cancel()
context.WithTimeout(parent, d)设置超时时间,到达时间后自动取消
context.WithDeadline(parent, t)设置绝对截止时间,到达时间后自动取消
context.WithValue(parent, key, val)为 Context 附加请求相关的键值数据
ctx.Done()返回一个 Channel,Context 被取消或超时时关闭
ctx.Err()返回 Context 结束的原因,通常是 context.Canceledcontext.DeadlineExceeded
ctx.Deadline()返回 Context 的截止时间
ctx.Value(key)读取通过 WithValue 写入的值

WithCancel 需要代码在合适的时机主动调用 cancel()WithTimeoutWithDeadline 则由定时器触发取消。取消信号只会从父 Context 向子 Context 传播,取消子 Context 不会取消父 Context。WithValue 适合传递请求范围内的少量数据,不应替代函数参数,也不应使用内置类型作为 key,避免不同包之间发生冲突。

在 HTTP 服务中,通常继续传递请求已有的 request.Context()。这样客户端断开连接时,取消信号能够沿调用链传播。

使用 Mutex 保护共享状态

Channel 适合在 Goroutine 之间传递数据,但多个 Goroutine 修改同一份状态时,也可以使用 sync.Mutex 保护临界区:

go
var mutex sync.Mutex
var group sync.WaitGroup

counter := 0

for i := 0; i < 100; i++ {
	group.Add(1)

	go func() {
		defer group.Done()

		mutex.Lock()
		defer mutex.Unlock()

		counter++
	}()
}

group.Wait()
fmt.Println(counter)

Lock() 保证同一时刻只有一个 Goroutine 修改 counterUnlock() 释放锁。如果移除锁,多个 Goroutine 可能同时读取并覆盖计数器,形成数据竞争。

可以使用竞态检测器检查程序运行期间发生的数据竞争:

bash
go run -race main.go
go test -race ./...

明确并发任务的生命周期

编写并发代码时,需要回答三个问题:

  1. 谁负责等待或取消 Goroutine;
  2. 谁拥有 Channel,并在必要时关闭它;
  3. 共享数据由谁保护。

Goroutine 的创建和切换成本较低,但仍会占用资源。不能退出的 Goroutine 会造成泄漏,错误的 Channel 操作可能导致 panic 或死锁,未受保护的共享状态则会产生数据竞争。

总结

Goroutine 用于并发执行函数,WaitGroup 用于等待任务完成,Channel 用于传递数据,Context 用于传递取消信号,Mutex 则用于保护共享状态。使用这些工具时,比启动任务更重要的是明确任务何时结束、数据由谁拥有以及共享状态如何同步。