Goroutine 与 Channel
这一章不是教你“会写 go func()”就结束,而是把 Go 并发当成生产系统能力来学:任务怎么启动,怎么通信,怎么限制并发,怎么取消,怎么收集错误,怎么避免泄漏,线上怎么排查。
学习目标
学完本章你应该能回答这些问题:
- goroutine 和系统线程是什么关系,为什么 goroutine 轻量但不能无限创建。
- channel 为什么既是通信工具,也是同步工具。
- 无缓冲、有缓冲、nil、已关闭 channel 分别会发生什么。
select、WaitGroup、context、errgroup在生产代码里分别解决什么问题。- worker pool、fan-out/fan-in、pipeline 怎么落地到采集、聚合、消息消费场景。
- goroutine 泄漏、数据竞争、阻塞、ticker 泄漏怎么定位。
并发不等于无限 goroutine
Go 让创建并发任务非常简单:
go func() {
fmt.Println("hello goroutine")
}()但这只是语法。生产上真正重要的是:这个 goroutine 什么时候退出,失败怎么返回,超时怎么停止,下游慢了会不会把本服务拖垮。
错误思路:
for _, task := range tasks {
go handle(task)
}如果 tasks 有 100 万条,这段代码会快速创建大量 goroutine。每个 goroutine 都需要栈、调度状态、可能还占用 HTTP 连接、数据库连接、内存缓冲区。结果可能是:
- 本服务内存上涨。
- Go 调度器压力变大。
- 下游接口被瞬间打爆。
- goroutine 阻塞后长期不退出。
- 日志里只能看到“服务慢”,但真正原因是并发失控。
正确思路是:并发必须有边界。
flowchart TD
A["任务很多"] --> B["进入有限队列"]
B --> C["固定数量 worker"]
C --> D["处理任务"]
D --> E["超时、失败、重试、退出"]goroutine 和线程
goroutine 是 Go 运行时管理的轻量任务,线程是操作系统调度的执行单元。Go 运行时会把很多 goroutine 调度到少量系统线程上执行。
flowchart TD
A["G: goroutine"] --> B["P: 调度上下文"]
B --> C["M: OS Thread"]
C --> D["CPU 执行"]| 对比点 | goroutine | OS 线程 |
|---|---|---|
| 管理者 | Go runtime | 操作系统 |
| 栈 | 初始栈小,可增长 | 通常固定且较大 |
| 创建成本 | 较低 | 较高 |
| 调度 | runtime 调度 | 内核调度 |
| 适合 | 大量 IO 并发、任务编排 | 底层执行单元 |
为什么 goroutine 轻量:
- 初始栈很小,不像系统线程通常需要更大的栈空间。
- Go runtime 可以在用户态完成一部分调度,减少频繁进入内核的成本。
- 阻塞在网络 IO、channel、锁等位置时,runtime 可以调度其他 goroutine 执行。
为什么仍然不能无限创建:
- goroutine 的栈再小也要占内存。
- goroutine 越多,调度、扫描、排查成本越高。
- 下游资源不是无限的,数据库连接池、HTTP 连接池、磁盘、CPU 都有上限。
goroutine 生命周期
一个生产级 goroutine 至少要能说清四件事:
- 谁创建它。
- 它依赖哪些输入。
- 它遇到错误怎么反馈。
- 它什么时候退出。
flowchart TD
A["创建 goroutine"] --> B{"等待任务或信号"}
B --> C["执行任务"]
C --> D{"完成、失败、取消、超时"}
D --> E["释放资源"]
E --> F["退出"]
B --> G["一直阻塞"]
G --> H["goroutine 泄漏"]容易泄漏的 goroutine 通常有这些特征:
- 阻塞在 channel 发送,但没人接收。
- 阻塞在 channel 接收,但没人发送也没人关闭。
- 死循环里没有退出条件。
- 外部 HTTP、数据库、RPC 没有设置超时。
- 请求结束后,后台 goroutine 没有监听
context。
channel 的本质
channel 是 goroutine 之间传递数据和同步执行时机的机制。
最小示例:
package main
import "fmt"
func main() {
ch := make(chan string)
go func() {
ch <- "hello"
}()
msg := <-ch
fmt.Println(msg)
}这里不是“发送方把值放进去就走了”这么简单。因为 ch 是无缓冲 channel,发送方和接收方必须同时准备好,数据才能交接。
flowchart TD
A["发送方执行 ch <- value"] --> B{"接收方是否已准备"}
B -- "否" --> C["发送方阻塞"]
B -- "是" --> D["值交给接收方"]
D --> E["双方继续执行"]无缓冲和有缓冲 channel
无缓冲 channel:
ch := make(chan int)特点:
- 发送和接收必须配对。
- 适合做同步交接。
- 能让生产者感知消费者是否跟得上。
有缓冲 channel:
ch := make(chan int, 10)特点:
- 缓冲未满时,发送可以先放入队列。
- 缓冲为空时,接收会阻塞。
- 缓冲满时,继续发送会阻塞。
- 适合做有限任务队列。
flowchart TD
A["发送数据"] --> B{"缓冲区满了吗"}
B -- "未满" --> C["写入缓冲区"]
B -- "已满" --> D["发送方阻塞"]
E["接收数据"] --> F{"缓冲区空了吗"}
F -- "未空" --> G["取出数据"]
F -- "已空" --> H["接收方阻塞"]不要把缓冲区设置得特别大来“解决慢”。大缓冲只是在掩盖消费者慢的问题,延迟会升高,内存会增长,失败恢复更困难。
channel 状态表
面试和生产排查里,channel 状态非常常问。
| channel 状态 | 发送 | 接收 | 关闭 |
|---|---|---|---|
| nil channel | 永久阻塞 | 永久阻塞 | panic |
| open 无缓冲 | 等接收方,可能阻塞 | 等发送方,可能阻塞 | 成功 |
| open 有缓冲未满 | 写入缓冲,通常不阻塞 | 有数据则返回,没数据阻塞 | 成功 |
| open 有缓冲已满 | 阻塞 | 取出数据 | 成功 |
| closed 且有缓冲数据 | panic | 先读缓冲数据,ok=true | 重复关闭 panic |
| closed 且无缓冲数据 | panic | 返回零值,ok=false | 重复关闭 panic |
示例:
package main
import "fmt"
func main() {
ch := make(chan int, 2)
ch <- 10
ch <- 20
close(ch)
for {
value, ok := <-ch
if !ok {
fmt.Println("channel closed")
break
}
fmt.Println(value)
}
}输出:
10
20
channel closed关键规则:
- 一般由发送方关闭 channel,因为发送方知道后续是否还会发送。
- 接收方不要关闭 channel,否则可能和发送方并发冲突。
- 不要重复关闭 channel。
- 不要向已关闭 channel 发送数据。
- nil channel 在
select里可以用来动态禁用某个 case,但直接发送或接收会永久阻塞。
select
select 用于同时等待多个 channel 操作。
select {
case msg := <-ch:
fmt.Println(msg)
case <-time.After(time.Second):
fmt.Println("timeout")
}执行规则:
- 如果多个 case 同时可执行,Go 会伪随机选择一个。
- 如果没有 case 可执行,又没有
default,当前 goroutine 阻塞。 - 如果没有 case 可执行,但有
default,立即执行default。 select常用于超时、取消、多个输入源、非阻塞尝试。
错误示例:default 导致 CPU 空转。
for {
select {
case msg := <-ch:
handle(msg)
default:
// 没有阻塞等待,会疯狂循环
}
}修复方式:
for {
select {
case msg := <-ch:
handle(msg)
case <-ctx.Done():
return
}
}如果确实要轮询,也要加 time.Ticker 或短暂 sleep,但更推荐用事件驱动。
WaitGroup
sync.WaitGroup 用于等待一组 goroutine 完成。
正确写法:
package main
import (
"fmt"
"sync"
)
func main() {
var wg sync.WaitGroup
for i := 0; i < 3; i++ {
wg.Add(1)
go func(i int) {
defer wg.Done()
fmt.Println(i)
}(i)
}
wg.Wait()
}核心规则:
Add要在启动 goroutine 之前调用。- goroutine 内部用
defer wg.Done(),避免 panic 或提前 return 导致计数不减。 WaitGroup使用后不要复制。Done次数不能超过Add次数,否则 panic。WaitGroup只能等待完成,不能自动收集错误,也不能自动取消其他 goroutine。
错误示例:在 goroutine 里面 Add。
for i := 0; i < 3; i++ {
go func() {
wg.Add(1) // 错误:Wait 可能已经看到计数为 0 并返回
defer wg.Done()
}()
}
wg.Wait()context 控制生命周期
context.Context 用于跨调用链传递取消信号、超时截止时间和少量请求级元数据。
flowchart TD
A["请求进入"] --> B["创建 ctx"]
B --> C["传给 service、db、http、goroutine"]
C --> D{"请求完成、超时、取消"}
D -- "是" --> E["ctx.Done 关闭"]
E --> F["下游停止工作并释放资源"]示例:
package main
import (
"context"
"fmt"
"time"
)
func watch(ctx context.Context) {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
fmt.Println("stop:", ctx.Err())
return
case <-ticker.C:
fmt.Println("tick")
}
}
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 2500*time.Millisecond)
defer cancel()
watch(ctx)
}为什么创建带 cancel 的 context 后要 defer cancel():
- 主动释放 timer 等资源。
- 主动通知下游 goroutine 退出。
- 避免任务已经结束但上下文资源还等到超时才清理。
worker pool:限制并发
worker pool 的目标是:任务可以很多,但同时执行的任务数有限。
商业场景:
- 医疗采集 Agent 批量读取接口数据。
- 批量调用三方 API 同步资产。
- MQ 消费者内部并发处理消息。
- 批量生成向量 Embedding。
生产级 worker pool 示例:
package main
import (
"context"
"errors"
"fmt"
"sync"
"time"
)
type Job struct {
ID int
}
func handle(ctx context.Context, job Job) error {
select {
case <-time.After(100 * time.Millisecond):
if job.ID == 7 {
return errors.New("mock downstream failed")
}
fmt.Println("handled", job.ID)
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func worker(ctx context.Context, id int, jobs <-chan Job, errorsCh chan<- error, wg *sync.WaitGroup) {
defer wg.Done()
for {
select {
case <-ctx.Done():
return
case job, ok := <-jobs:
if !ok {
return
}
if err := handle(ctx, job); err != nil {
select {
case errorsCh <- fmt.Errorf("worker %d job %d: %w", id, job.ID, err):
case <-ctx.Done():
return
}
}
}
}
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
jobs := make(chan Job, 10)
errorsCh := make(chan error, 20)
var wg sync.WaitGroup
workerCount := 4
for i := 0; i < workerCount; i++ {
wg.Add(1)
go worker(ctx, i, jobs, errorsCh, &wg)
}
for i := 0; i < 20; i++ {
jobs <- Job{ID: i}
}
close(jobs)
wg.Wait()
close(errorsCh)
for err := range errorsCh {
fmt.Println("error:", err)
}
}这个 Demo 的关键点:
workerCount控制并发上限。jobs是有限缓冲,不让任务无限堆在内存里。context负责超时和取消。errorsCh收集失败,不让错误悄悄丢失。- 发送方关闭
jobs,等待所有 worker 退出后关闭errorsCh。
errgroup:带错误传播的并发
WaitGroup 只能等待,不能返回错误。生产中经常希望:只要一个子任务失败,就取消其他任务,并把错误返回给上层。可以使用 golang.org/x/sync/errgroup。
安装:
go get golang.org/x/sync/errgroup示例:接口聚合。
package main
import (
"context"
"fmt"
"time"
"golang.org/x/sync/errgroup"
)
func query(ctx context.Context, name string) error {
select {
case <-time.After(200 * time.Millisecond):
fmt.Println("done:", name)
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
group, ctx := errgroup.WithContext(ctx)
for _, name := range []string{"user", "order", "asset"} {
name := name
group.Go(func() error {
return query(ctx, name)
})
}
if err := group.Wait(); err != nil {
fmt.Println("failed:", err)
return
}
fmt.Println("all success")
}errgroup 适合:
- 多个外部接口并发查询。
- 多个文件并发解析。
- 多个任务任一失败就停止整批任务。
- 需要错误返回给调用方的并发编排。
不适合:
- 需要持续消费的长期 worker。
- 需要收集所有错误而不是遇到第一个错误就返回的场景。
- 复杂重试、限流、队列调度,此时需要专门设计任务模型。
fan-out 和 fan-in
fan-out 是把任务分发给多个 goroutine 并发处理,fan-in 是把多个结果合并回来。
flowchart TD
A["输入任务流"] --> B["worker 1"]
A --> C["worker 2"]
A --> D["worker 3"]
B --> E["结果合并"]
C --> E
D --> E示例:批量计算。
func worker(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for {
select {
case <-ctx.Done():
return
case value, ok := <-in:
if !ok {
return
}
select {
case out <- value * value:
case <-ctx.Done():
return
}
}
}
}()
return out
}生产要点:
- 所有阶段都要监听
ctx.Done()。 - 下游提前退出时,上游不能永久阻塞在发送。
- 合并结果时要等待所有 worker 关闭输出。
- 任务量大时,要限制 worker 数。
pipeline:分阶段处理
采集、清洗、校验、入库适合用 pipeline 思路。
flowchart TD
A["读取数据"] --> B["解析"]
B --> C["校验"]
C --> D["脱敏"]
D --> E["批量入库或上报"]简单版 Demo:
func read(ids []int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, id := range ids {
out <- id
}
}()
return out
}
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for value := range in {
out <- value * value
}
}()
return out
}这个 Demo 只适合理解模型。生产版本必须加入:
- context 取消。
- 错误通道。
- 有限缓冲。
- 下游慢时的背压。
- 每个阶段的指标和日志。
Timer 和 Ticker
定时相关并发也很容易出问题。
常见工具:
| 工具 | 用途 | 注意点 |
|---|---|---|
time.After | 等待一次超时 | 高频循环里可能制造很多 timer |
time.NewTimer | 可停止的一次性 timer | 不用时调用 Stop |
time.NewTicker | 周期性触发 | 不用时调用 Stop |
time.Tick | 简写 ticker | 长生命周期代码中不推荐,无法主动 Stop |
推荐写法:
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
doWork()
}
}不这样做会怎样:
- 定时器长期存在,资源不释放。
- goroutine 退出了,但 ticker 仍然尝试触发。
- 服务运行时间越久,隐藏问题越明显。
数据竞争
数据竞争是:多个 goroutine 同时访问同一变量,其中至少一个是写,并且没有同步保护。
错误示例:
count := 0
var wg sync.WaitGroup
for i := 0; i < 1000; i++ {
wg.Add(1)
go func() {
defer wg.Done()
count++
}()
}
wg.Wait()
fmt.Println(count)count++ 不是原子操作,它至少包含读、加、写三个步骤。多个 goroutine 同时执行会互相覆盖。
修复方式一:Mutex。
var mu sync.Mutex
mu.Lock()
count++
mu.Unlock()修复方式二:atomic。
var count int64
atomic.AddInt64(&count, 1)检查方式:
go test -race ./...选择建议:
| 工具 | 适合 | 不适合 |
|---|---|---|
| Mutex | 保护 map、结构体、缓存等共享状态 | 表达复杂任务流 |
| RWMutex | 读多写少 | 写频繁 |
| atomic | 简单计数、状态位 | 复杂结构 |
| channel | 任务传递、事件通知、退出信号 | 简单字段自增 |
| sync.Map | 特定高并发读写 map | 普通业务 map 滥用 |
goroutine 泄漏排查
判断流程:
flowchart TD
A["发现内存或 goroutine 数上涨"] --> B["查看 runtime.NumGoroutine 或监控"]
B --> C["抓 pprof goroutine"]
C --> D{"堆栈卡在哪里"}
D --> E["chan send: 没人接收"]
D --> F["chan receive: 没人发送或关闭"]
D --> G["Mutex.Lock: 锁竞争或死锁"]
D --> H["net/http: 下游慢或无超时"]
D --> I["database/sql: 连接池等待"]
E --> J["补接收、缓冲、取消或关闭"]
F --> J
G --> K["缩小锁范围或修复死锁"]
H --> L["补超时、限流、熔断"]
I --> M["查慢 SQL、连接池、事务未释放"]开启 pprof:
import _ "net/http/pprof"
go func() {
_ = http.ListenAndServe(":6060", nil)
}()查看:
go tool pprof http://127.0.0.1:6060/debug/pprof/goroutine直接看文本堆栈:
http://127.0.0.1:6060/debug/pprof/goroutine?debug=2代码里也可以临时观察:
fmt.Println("goroutines:", runtime.NumGoroutine())常见堆栈含义:
| 堆栈关键词 | 可能原因 |
|---|---|
chan send | 发送没人接收 |
chan receive | 接收没人发送,或 channel 没关闭 |
select | 等待多个信号,可能缺少取消路径 |
sync.Mutex.Lock | 锁竞争、死锁、锁范围过大 |
net/http | 外部接口慢、连接未释放、没设置 timeout |
database/sql | 连接池耗尽、慢 SQL、事务没提交或回滚 |
商业场景:医疗采集 Agent
采集 Agent 常见链路:
flowchart TD
A["读取医院接口、文件、设备数据"] --> B["解析结构"]
B --> C["校验字段和编码"]
C --> D["脱敏和标准化"]
D --> E["批量上报平台"]
E --> F["失败落盘和重试"]并发设计:
- 读取端限速,避免把医院接口打挂。
- 解析阶段可以用 worker pool。
- 上报阶段必须有 HTTP timeout、重试上限和幂等 key。
- context 控制程序停止,不能强杀导致数据丢失。
- 失败数据本地落盘,重启后继续补偿。
- 指标记录队列长度、处理耗时、失败数、goroutine 数。
商业场景:接口聚合
例如一个资产详情页需要并发查询用户、资产、订单、权限、统计。
核心原则:
- 总超时由入口 context 控制。
- 单个下游接口也要有自己的 timeout。
- 强依赖失败则整体失败,弱依赖失败可以降级。
- 使用
errgroup管理错误和取消。 - 日志里带 trace_id 和每个下游耗时。
商业场景:MQ 消费者并发
MQ 消费者不是并发越大越好。
要控制:
- 消费 goroutine 数。
- 每条消息处理超时。
- 下游数据库连接池。
- 幂等处理。
- 失败重试和死信。
- 退出时先停止拉取,再等待正在处理的消息完成。
如果无限开 goroutine 处理消息,短时间看消费速度变快,后面可能因为数据库连接池耗尽、锁竞争、下游限流导致堆积更严重。
面试标准回答
goroutine 为什么轻量
goroutine 是 Go runtime 管理的轻量任务,初始栈小且可增长,runtime 通过 GMP 模型把大量 goroutine 调度到较少的系统线程上执行,所以创建和切换成本低于直接创建大量 OS 线程。但 goroutine 不是免费资源,生产上仍然要限制数量、控制生命周期,并处理超时、取消和泄漏。channel 阻塞规则
无缓冲 channel 发送和接收必须同时准备好,否则会阻塞。有缓冲 channel 在缓冲未满时发送不阻塞,缓冲为空时接收阻塞,缓冲满时发送阻塞。nil channel 发送和接收会永久阻塞,关闭 nil channel 会 panic。已关闭 channel 可以继续接收缓冲数据,读完后返回零值和 ok=false,但不能再发送。WaitGroup 和 errgroup 区别
WaitGroup 只能等待一组 goroutine 结束,不负责返回错误,也不会自动取消其他 goroutine。errgroup 在 WaitGroup 的基础上支持错误返回,并且可以结合 context 在任一 goroutine 出错时取消其他任务,适合接口聚合、批处理等需要错误传播的场景。goroutine 泄漏怎么排查
先看 goroutine 数是否持续增长,再抓 pprof goroutine profile,看堆栈卡在 channel 发送、channel 接收、锁、网络 IO、数据库连接还是定时器上。然后回到代码检查是否监听 context,channel 是否由发送方关闭,下游调用是否有超时,循环是否有退出条件。Mutex、channel、atomic 怎么选
Mutex 适合保护共享状态,channel 适合任务传递和事件通知,atomic 适合简单计数和状态标记。不能为了追求所谓 Go 风格强行用 channel 替代清晰的锁,也不能在需要任务流和取消信号的场景里到处共享变量。关联知识点
本章小结
Go 并发的难点不在语法,而在工程边界。生产代码必须能回答:goroutine 什么时候退出,channel 谁关闭,并发数谁限制,错误怎么返回,超时怎么取消,泄漏怎么排查。答不上这些问题,并发代码越多,系统越不稳定。
