大家好,我是小蒙。
刚开始写 Go 的时候,谁还没干过「来一个请求起一个 go func()」的事?当时觉得并发真爽,直到后来线上流量一冲,GC 压力拉满、内存爆掉、协程泄漏把系统直接打垮,我才老实。
今天不整虚的八股文,直接聊聊我在实际业务里最常压测和落地的 4 种 Go 并发模式,附带关键代码片段和 Benchmark 对比数据。
1. Worker Pool(工作池模式):告别无节制的 Goroutine
业务场景
批量导入 Excel 数据、消费大量消息队列任务。如果直接对每条数据起 Goroutine,面对上百万数据瞬间就 OOM 了。
核心代码
package main
import (
"context"
"sync"
)
type Task struct {
ID int
}
func Worker(ctx context.Context, id int, tasks <-chan Task, wg *sync.WaitGroup) {
defer wg.Done()
for {
select {
case <-ctx.Done():
return
case task, ok := <-tasks:
if !ok {
return
}
// 处理业务逻辑
_ = task.ID * 2
}
}
}
func RunWorkerPool(ctx context.Context, workerCount int, taskList []Task) {
tasks := make(chan Task, 100)
var wg sync.WaitGroup
for i := 0; i < workerCount; i++ {
wg.Add(1)
go Worker(ctx, i, tasks, &wg)
}
for _, task := range taskList {
tasks <- task
}
close(tasks)
wg.Wait()
}
Benchmark 实测
处理 100,000 个任务(模拟轻量计算 + 内存开销),对比无脑并发 vs 固定 50 个 Worker 的 Pool:
goos: darwin
goarch: arm64
pkg: demo/concurrency
Benchmark_UnboundedGoroutines-8 12 92410200 ns/op 24510200 B/op 100120 allocs/op
Benchmark_WorkerPool_50-8 48 24150300 ns/op 802340 B/op 124 allocs/op
点评:耗时缩短了约 73%,内存分配从 24 MB 骤降到 800 KB。不仅 GC 压力归零,调度开销也小了极多。
2. Fan-Out / Fan-In(扇出 / 扇入模式):微服务数据聚合必备
业务场景
B 端首页加载,需要并发调用「用户画像服务」、「订单中心」、「风控系统」,最后汇总为一个 VO。
核心思路
- Fan-Out:主协程派发多个子协程独立去拉取不同 RPC / HTTP 接口。
- Fan-In:用一个统一的 Channel 或
sync.WaitGroup收集结果,搭配context.WithTimeout防止被慢接口拖垮。
核心代码
func FetchAggregatedData(ctx context.Context) ([]string, error) {
ctx, cancel := context.WithTimeout(ctx, 300*time.Millisecond)
defer cancel()
results := make(chan string, 3)
var wg sync.WaitGroup
sources := []func(context.Context) string{
fetchUserProfile,
fetchOrderStats,
fetchRiskScore,
}
// Fan-Out
for _, fn := range sources {
wg.Add(1)
go func(f func(context.Context) string) {
defer wg.Done()
select {
case <-ctx.Done():
return
case results <- f(ctx):
}
}(fn)
}
// 监听完成并关闭 Channel
go func() {
wg.Wait()
close(results)
}()
// Fan-In 收集
var finalData []string
for res := range results {
finalData = append(finalData, res)
}
if ctx.Err() != nil {
return nil, ctx.Err()
}
return finalData, nil
}
小蒙踩坑提醒:子协程往 Channel 发数据时,一定要
select + ctx.Done(),否则如果外部超时退出导致没有 Receiver,子协程会永久卡死在results <- res,造成 Goroutine 泄漏!
3. Pipeline(管道模式):流式清洗大数据
业务场景
解析 Nginx 访问日志:读取原始日志 -> 正则清洗 -> IP 归属地转换 -> 入库 ClickHouse。
核心模式
每个阶段(Stage)都是一个独立的 Goroutine,通过 In / Out Channel 串联起来,流式处理,内存占用恒定。
// Stage 1: 生成数据
func gen(ctx context.Context, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case <-ctx.Done():
return
case out <- n:
}
}
}()
return out
}
// Stage 2: 处理数据 (平方)
func sq(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
select {
case <-ctx.Done():
return
case out <- n * n:
}
}
}()
return out
}
这种写法的优雅之处在于:哪怕处理 100 GB 的日志文件,内存占用也可以死死压在几十兆以内。
4. Singleflight:防击穿合并请求
业务场景
高并发热点 Key 失效瞬间(比如爆款商品秒杀),数万个请求同时穿透到 MySQL。
解决方案
利用 Go 官方拓展包 golang.org/x/sync/singleflight,让同一时刻成百上千个相同的查询合并为只有 1 个去查库,其余并发请求原地等待该结果共享。
核心代码
import "golang.org/x/sync/singleflight"
var g singleflight.Group
func GetProductDetail(productID string) (string, error) {
// 命中本地或 Redis 缓存省略...
// 缓存未命中,使用 Singleflight 合并下游 DB 请求
v, err, shared := g.Do(productID, func() (interface{}, error) {
// 模拟慢查询 DB
time.Sleep(50 * time.Millisecond)
return queryDB(productID)
})
if err != nil {
return "", err
}
_ = shared // shared == true 表示该请求共享了其他人的结果
return v.(string), nil
}
压测对比(5000 QPS 并发打同一热点 Key)
- 未使用 Singleflight:DB 瞬间承压 5000 QPS,连接池被打满,P99 延迟飙到 800ms+。
- 使用 Singleflight:DB 实际 QPS 降为 20 ~ 30 QPS,P99 稳定在 52ms 左右。
总结
写 Go 并发千万不要盲目追求「起协程快」。评判并发写得好不好的标准其实就三点:
- 生命周期是否可控?(有没有泄漏?Context 是否传透了?)
- 并发量是否有上限?(有没有池化或者令牌桶控频?)
- 分配开销是否划算?(Channel 和 锁 谁更合适?用不用
sync.Pool?)
大家在生产环境最常用哪种模式?遇到过什么诡异的并发 Bug?欢迎在评论区一起交流,我也常去扒 pprof 看调度日志。
许可协议:CC BY-NC 4.0
更新于 2 小时前
觉得文章有帮助?点个赞吧!
0 条评论


