Runtime 调度器 — G(协程) / P(逻辑处理器) / M(内核线程) 三层抽象: work stealing · syscall hand-off · 信号抢占
GMP 把"要执行的代码 (G)"、"执行凭证 (P)"、"系统线程 (M)"解耦: goroutine 便宜到可以随手百万个, P 数量(GOMAXPROCS)决定真并行度, M 只是在 P 上干活的内核线程。三大机制 — work stealing / syscall hand-off / 信号抢占 — 保证 CPU 几乎永不空闲。
for i := 0; i < 1000000; i++ { go handle(reqs[i]) // 关键: 创建纯用户态 ns 级, 栈 2KB 起步 } // → runtime.NumGoroutine() ≈ 100 万, OS 线程仍只有几十个
GOMAXPROCS = 真并行度上限, 默认 CPU 核数。runtime.GOMAXPROCS(0) // → 10 (只查询, 不修改) runtime.GOMAXPROCS(2) // 真并行度上限压到 2 // 关键: 20 个 CPU 密集 G 也最多同时占 2 个核
# 大量阻塞 syscall 时观察线程数: GODEBUG=schedtrace=1000 ./app # → SCHED ... gomaxprocs=8 idleprocs=5 threads=157 ... # 关键: threads(=M) 157 远大于 gomaxprocs(=P) 8
ch := make(chan int) go func() { ch <- 1 }() // 接收方唤醒后进其 P 的 runnext v := <-ch // → 1, 下一调度周期立即被挑中
// runtime findRunnable 示意 (proc.go): if schedtick%61 == 0 && runqsize > 0 { globrunqget(pp, 1) // 关键: 每 61 次必摸全局队列 }
# 每个 P 各压满任务, 观察偷取是否发生: GODEBUG=schedtrace=1000 ./app # → SCHED ... gomaxprocs=8 idleprocs=0 ... # 关键: 均衡全自动, 代码里不写任何分发逻辑
SIGURG 信号强制抢占。go func() { for {} }() // 定长死循环 go func() { fmt.Println("alive") }() // 1.14+: → 正常打印 // 1.13-: 同 P 的第二个 G 被饿死 (无调用点不让出)
v := <-ch // ① chan: G 入等待队列, M 换下一个 G n, _ := conn.Read(buf) // ② 网络: G 入 netpoller (epoll) f.Sync() // ③ syscall: M 陷内核, P hand-off // 关键: 三种阻塞都不占用稀缺的 P
net/http 默认就是"每连接两 goroutine"模型, 这是 Go 单机轻松十万级并发的根基:
http.HandleFunc("/order", func(w http.ResponseWriter, r *http.Request) { go audit.Log(r) // 异步审计: 随手一个 go, 2KB 起步, 不心疼 result := query(r) // 阻塞 IO 期间调度器自动让 M 去服务其他请求 json.NewEncoder(w).Encode(result) })
K8s 给 Pod 限 CPU=2, 但宿主机 64 核时旧版 Go 默认 GOMAXPROCS=64: 64 个 P 争抢 2 核配额, 被 CFS throttle 后出现周期性延迟毛刺(P99 飙升)。业界标准解法(Uber 方案):
import _ "github.com/uber-go/automaxprocs" // 自动读 cgroup quota, 设 GOMAXPROCS=配额核数 // 或 main 里显式: debug.SetGCPercent / runtime.GOMAXPROCS(2)
批量调下游接口, 并发不能超过 50(下游限流), 用 buffered channel 当信号量:
sem := make(chan struct{}, 50) // 容量 50 = 最多 50 个并发 wg := &sync.WaitGroup{} for _, id := range ids { wg.Add(1) sem <- struct{}{} // 满了就阻塞在这里, 等价于排队 go func(id string) { defer wg.Done(); defer func() { <-sem }() fetch(id) }(id) } wg.Wait()
并行拉多个下游、任一失败全体停止 —— errgroup 是现代标准写法:
g, ctx := errgroup.WithContext(ctx) g.SetLimit(10) // 内置并发上限 (1.20+) for _, id := range ids { g.Go(func() error { return fetch(ctx, id) }) } err := g.Wait() // 任一 error → ctx 取消 → 其余快速退出
指标单调上涨 = 泄漏; pprof 按创建栈聚合, 找卡在 chan/lock 的栈:
// 1. 暴露指标: goroutine 数趋势 (告警看斜率, 不看绝对值) prometheus.MustRegister(collectors.NewGoCollector()) log.Info("goroutines", "n", runtime.NumGoroutine()) # 2. 线上抓现场: 按创建位置聚合, 找最多的卡点栈 curl localhost:6060/debug/pprof/goroutine?debug=1 | head -50 # 3. 常见结论: "chan send" 无接收 / "select 无 ctx.Done 分支" → 补取消
埋点上报: 累积到量或到时即刷, 停机强制 flush:
func (b *Buffer) run(ctx context.Context) { ticker := time.NewTicker(2 * time.Second) defer ticker.Stop() for { select { case <-ticker.C: // 到时刷 b.flush() case e := <-b.in: b.buf = append(b.buf, e) if len(b.buf) >= 1000 { // 到量提前刷 b.flush() } case <-ctx.Done(): // 停机: 强制刷尾批再退 b.flush() return } } }
多 GB 文件按偏移分段, 每段一 goroutine 用 ReadAt 并发读(句柄并发安全):
func parallelSum(f *os.File, size int64) uint64 { n := runtime.GOMAXPROCS(0) // 分段数 = P 的 1~2 倍 seg := size / int64(n) results := make(chan uint64, n) for i := 0; i < n; i++ { go func(off, length int64) { buf := make([]byte, length) f.ReadAt(buf, off) // ReadAt 并发安全 results <- checksum(buf) }(int64(i)*seg, seg) } var total uint64 for i := 0; i < n; i++ { total += <-results } return total }
调度健康度三件套: goroutine 数、GC 停顿、GOMAXPROCS, 客户端库自带:
import _ "github.com/prometheus/client_golang/prometheus/auto" # 已自动暴露: go_goroutines / go_gc_duration_seconds / go_memstats_* # 告警规则示例 (PromQL): # deriv(go_goroutines[10m]) > 0.5 → goroutine 持续增长疑似泄漏 # rate(go_gc_duration_seconds_sum[5m]) > 0.1 → GC 吃掉 10% 时间
每连接两 goroutine(读+写), done chan 双向退出 —— 单机十万连接的常规架构:
func (s *Server) serveConn(conn net.Conn) { done := make(chan struct{}) outbox := make(chan []byte, 64) go s.writeLoop(conn, outbox, done) // 唯一写者: 不交错 go s.heartbeat(conn, outbox, done) // 定期 ping 保活 s.readLoop(conn, outbox) // 阻塞读, 出错即返回 close(done) // 广播退出: 写/心跳循环感知并收尾 }
行情/风控类服务把 P 与核绑定, 减少跨核迁移的缓存失效:
// 容器/K8s: cpuset 固定 2-4 号核, Go 自动感知 GOMAXPROCS=核数 docker run --cpuset-cpus=2-4 app # 裸机 systemd: CPUAffinity=2 3 4 # 验证: 服务内打点 runtime.GOMAXPROCS(0) 应等于绑定的核数 # 效果: P 不跨核迁移 → 缓存/TLB 命中率升 → P99 尾部更稳
// 错: for _, r := range reqs { go handle(r) } → 瞬时 10 万 G sem := make(chan struct{}, 100) // 对: 在途上限 100 for _, r := range reqs { sem <- struct{}{} // 满了排队, 不打爆 go func(r Req) { defer func() { <-sem }(); handle(r) }(r) }
context 超时取消; 上线后盯 pprof 的 goroutine 数是否单调上涨。// 错: ch <- result (无人接收, G 永久卡在 send) select { // 对: 一律带超时/取消 case ch <- result: case <-ctx.Done(): // → context canceled, G 正常退出 }
# 错: CPU limit=2 但宿主机 64 核 → GOMAXPROCS=64, 被 CFS throttle # 对: 自动对齐 cgroup 配额 import _ "github.com/uber-go/automaxprocs" # → GOMAXPROCS=2
for _, v := range 的 v 是复用变量, 闭包捕获的是同一个地址, 并发下全拿到最后一个值。Go 1.22 起每轮迭代新变量, 已修复; 老代码仍要 v := v。for _, v := range items { go func() { use(v) }() // 错(1.21-): 并发下全拿到最后一个 v go func(v int) { use(v) }(v) // 对: 按值传参 (或每轮 v := v) }
// 错: go do(); time.Sleep(time.Second) → 时序碰运气 done := make(chan struct{}) // 对: 等"什么"就等"什么" go func() { do(); close(done) }() <-done // → 精确同步, 零竞态
defer func(){ if r := recover(); r != nil { log } }(), 或统一封装 go safe(fn)。go func() { panic("boom") }() // 错: 裸 go, 整个进程被打死 // 对: 统一封装, 入口顶层 recover func safe(fn func() error) { defer func() { if r := recover(); r != nil { log.Error(r) } }() fn() }
for _, f := range files { fh, _ := os.Open(f) defer fh.Close() // 错: 万个句柄攒到函数返回才释放 } for _, f := range files { process(f) } // 对: defer 进 process 内随轮归还
for _, t := range tasks { go func() { wg.Add(1); t() }() // 错: Wait 抢先返回, 少算 wg.Add(1) // 对: Add 在 go 语句之前 }
fatal error: all goroutines are asleep。正解: 无缓冲 chan 的 send/recv 必须分居两方; 单流程用有缓冲或变量直传。ch := make(chan int) // 错: 无缓冲 ch <- 1; fmt.Println(<-ch) // → fatal error: all goroutines are asleep ch2 := make(chan int, 1) // 对: 有缓冲 (或 send 放进另一个 G) ch2 <- 1; fmt.Println(<-ch2) // → 1
for { select { case <-ctx.Done(): // 对: 第一分支固定写退出 return case job := <-jobs: // 只有此分支 = 错形态, 永转不停 run(job) } }
for i := 0; i < 1e9; i++ { runtime.Gosched() // 错(过时): 1.14+ 异步抢占已内置 work(i) } for i := 0; i < 1e9; i++ { work(i) } // 对: 直接算; 超长任务下沉 worker
// 错: go func() { audit(order) }() → 日志无 traceId, 排查断链 go func(ctx context.Context, o Order) { // 对: ctx 作为第一个参数 log.From(ctx).Info("audit") // → traceId 自动注入 audit(ctx, o) }(ctx, order)
# 错: 无监控, 上线三天 OOM 才发现泄漏 # 对: go_goroutines 告警看斜率, 不看绝对值 deriv(go_goroutines[10m]) > 0.5 # → 持续增长即疑似泄漏
for !done { time.Sleep(...) } 是竞态+浪费的双重反模式。正解: done chan / WaitGroup / errgroup, 让"等待"语义显式。for !done { time.Sleep(100 * time.Millisecond) } // 错: 竞态+空转 done := make(chan struct{}) // 对: 事件驱动 close(done); <-done // → 精确唤醒
func main() { go flusher() // 错: main 一退 flusher 立即死 serve() // → 尾批数据/未 ack 消息丢失 wg.Wait(); flush() // 对: 退出前 Wait + 强制 flush }
big := make([]byte, 64<<20) // 64MB go func() { use(big) }() // 错: 64MB 被 G 拉满整个生命周期 go func(n int) { calc(n) }(len(big)) // 对: 只带需要的值
defer log.Println("defer 仍执行") // Goexit 与 panic 一样跑 defer runtime.Goexit() // 只终止当前 G, 调用方无感知 // 错: 测试里用 Goexit 跳过断言 → 静默假绿 // 对: 断言失败直接 t.Fatal("bad state")
go func() { // 错: 内层新 G 不带 ctx go cleanup() // 外层取消它也不知道, 泄漏转移 }() go func() { go cleanup(ctx) }() // 对: 派生一律用传入的 ctx
# 错: 本机 16 核压测 P99=10ms, 直接外推生产 # 对: 与生产同限核同参数再压 GOMAXPROCS=2 ./bench # 模拟容器 CPU limit=2 docker run --cpus=2 --memory=4g app # → 压出真实容量
for _, job := range jobs { go handle(job) // 错: 无上限无生命周期, 排队靠运气 } g, ctx := errgroup.WithContext(ctx) // 对: 限流+取消+等待 g.SetLimit(50) // → 上限 50, 可取消, 可等待 for _, job := range jobs { g.Go(func() error { return handle(ctx, job) }) } g.Wait()