Goroutine、Channel 与并发控制
Go 通过 Goroutine 执行并发任务,并使用 Channel、Context 和 sync 包中的工具协调这些任务。本文从启动一个 Goroutine 开始,逐步介绍任务等待、数据传递、取消和共享状态保护。
启动 Goroutine
普通函数调用会按照当前控制流同步执行。在调用前添加 go,可以让函数在新的 Goroutine 中执行:
func printMessage() {
fmt.Println("hello")
}
func main() {
go printMessage()
}并发表示多个任务可以交错推进,并不保证它们在同一时刻并行运行。是否真正并行还取决于 CPU 核心数和 Go 运行时的调度。
上面的程序也不能保证输出 hello。如果 main() 先结束,整个进程会直接退出,不会继续等待其他 Goroutine。启动并发任务时,还需要明确由谁等待和结束它。
使用 WaitGroup 等待任务
只需要等待一组任务全部完成时,可以使用 sync.WaitGroup:
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:
messages := make(chan string)
go func() {
messages <- "done"
}()
message := <-messages
fmt.Println(message)有缓冲 Channel 可以暂存一定数量的值。缓冲区写满时,新的发送操作会等待;缓冲区为空时,接收操作会等待:
messages := make(chan string, 2)
messages <- "first"
messages <- "second"
fmt.Println(<-messages) // first
fmt.Println(<-messages) // secondChannel 按值成功发送的先后顺序进行接收,可以理解为先进先出(FIFO)。但多个 Goroutine 的执行顺序并不确定,因此并发发送时,哪个值先进入 Channel 仍取决于实际调度。
函数参数还可以限制 Channel 的操作方向:
func produce(output chan<- string) {
output <- "done"
}
func consume(input <-chan string) {
fmt.Println(<-input)
}chan<- string 只能发送,<-chan string 只能接收。方向限制可以让函数的职责更清楚。
使用 Channel 构建任务队列
多个 Goroutine 可以从同一个 Channel 接收任务。下面的程序启动两个 Worker,处理五个任务:
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()可以按以下步骤理解这段代码:
jobs := make(chan int, 3)创建容量为3的任务队列。缓冲区未满时,发送任务不会阻塞;缓冲区满后,发送方会等待 Worker 取走任务。for workerID := 1; workerID <= 2; workerID++启动两个 Worker。每个 Worker 用for jobID := range jobs从同一个 Channel 接收任务,直到jobs关闭并被读空。group.Add(1)在启动 Worker 前增加等待计数,defer group.Done()保证 Worker 退出时减少计数,group.Wait()等待两个 Worker 都结束。- 发送完 5 个任务后,发送方调用
close(jobs)。Worker 会先读空剩余任务,再结束range循环;如果不关闭 Channel,range会一直等待新任务。
Channel 的关闭
只有发送方确定不会再写入数据,并且接收方需要感知结束时,才需要关闭 Channel。重复关闭 Channel 或向已关闭的 Channel 发送数据都会触发 panic。
使用 select 等待多个操作
select 可以同时等待多个 Channel 操作,并执行最先就绪的分支,并略过其余的分支:
select {
case message := <-messages:
fmt.Println(message)
case <-time.After(time.Second):
fmt.Println("timeout")
}如果一秒内没有收到消息,time.After() 返回的 Channel 会就绪,程序进入超时分支。
使用 Context 传递取消信号
仅让等待方超时还不够,正在执行的 Goroutine 也应有机会结束。context.Context 可以在调用链中传递取消信号和截止时间:
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.Canceled 或 context.DeadlineExceeded |
ctx.Deadline() | 返回 Context 的截止时间 |
ctx.Value(key) | 读取通过 WithValue 写入的值 |
WithCancel 需要代码在合适的时机主动调用 cancel();WithTimeout 和 WithDeadline 则由定时器触发取消。取消信号只会从父 Context 向子 Context 传播,取消子 Context 不会取消父 Context。WithValue 适合传递请求范围内的少量数据,不应替代函数参数,也不应使用内置类型作为 key,避免不同包之间发生冲突。
在 HTTP 服务中,通常继续传递请求已有的 request.Context()。这样客户端断开连接时,取消信号能够沿调用链传播。
使用 Mutex 保护共享状态
Channel 适合在 Goroutine 之间传递数据,但多个 Goroutine 修改同一份状态时,也可以使用 sync.Mutex 保护临界区:
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 修改 counter,Unlock() 释放锁。如果移除锁,多个 Goroutine 可能同时读取并覆盖计数器,形成数据竞争。
可以使用竞态检测器检查程序运行期间发生的数据竞争:
go run -race main.go
go test -race ./...明确并发任务的生命周期
编写并发代码时,需要回答三个问题:
- 谁负责等待或取消 Goroutine;
- 谁拥有 Channel,并在必要时关闭它;
- 共享数据由谁保护。
Goroutine 的创建和切换成本较低,但仍会占用资源。不能退出的 Goroutine 会造成泄漏,错误的 Channel 操作可能导致 panic 或死锁,未受保护的共享状态则会产生数据竞争。
总结
Goroutine 用于并发执行函数,WaitGroup 用于等待任务完成,Channel 用于传递数据,Context 用于传递取消信号,Mutex 则用于保护共享状态。使用这些工具时,比启动任务更重要的是明确任务何时结束、数据由谁拥有以及共享状态如何同步。