分析一个简单的goroutine资源池

admin 2026-08-14 08:15:21 网络安全文章 来源:ZONE.CI 全球网 0 阅读模式

文章总结: 针对Go无节制创建goroutine导致线上卡死的问题,本文指出需建立并发边界。作者基于channel实现轻量级资源池,强调三处关键设计:任务提交须快速失败或设超时防阻塞,worker内须recover防脏数据打挂,关闭时须先close再Wait。建议按任务特性配置worker数,通过池化约束并发以保障系统稳定。 综合评分: 91 文章分类: 其他


cover_image

分析一个简单的goroutine资源池

原创

go go

Go语言教程

2026年7月19日 13:20 陕西

在小说阅读器读本章

去阅读

线上导入任务卡住,接口没报错,CPU 也没打满,日志停在一半。

这种问题我一般不先翻业务代码,先看 goroutine。

一看就很扎眼:

goroutines: 18432
heap: 910MB
last log: submit bill task, file=202606.csv, row=43891

这种数量基本不用猜,八成是哪里把 goroutine 当免费线程用了。

Go 写并发太顺手,顺手到容易出事。很多人一看到批量处理,就这么写:

for _, row := range rows {
    go syncBill(row)
}

开发环境跑 100 行数据,没问题。线上一批 10 万行,直接把下游接口、数据库连接池、机器内存一起拖下水。

goroutine 便宜,不等于不要钱。

我更愿意把它当成“要管理的资源”。既然是资源,就得有上限,有队列,有退出,有兜底。最简单的办法,就是搞一个 goroutine 池。

别一上来就找复杂框架,很多场景一个 channel 就够了。

我一般会先写成这样:

package pool

import (
 "context"
 "fmt"
 "log"
 "sync"
 "time"
)

type Job struct {
 ID      string
 Payload string
}

type WorkerPool struct {
 size int
 jobs chan Job
 wg   sync.WaitGroup
}

func NewWorkerPool(size int, queueSize int) *WorkerPool {
&nbsp;if&nbsp;size <=&nbsp;0&nbsp;{
&nbsp; size =&nbsp;1
&nbsp;}
&nbsp;if&nbsp;queueSize <=&nbsp;0&nbsp;{
&nbsp; queueSize = size *&nbsp;2
&nbsp;}

&nbsp;return&nbsp;&WorkerPool{
&nbsp; size: size,
&nbsp; jobs:&nbsp;make(chan&nbsp;Job, queueSize),
&nbsp;}
}

func&nbsp;(p *WorkerPool)&nbsp;Start(ctx context.Context)&nbsp;{
&nbsp;for&nbsp;i :=&nbsp;0; i < p.size; i++ {
&nbsp; workerID := i +&nbsp;1
&nbsp; p.wg.Add(1)

&nbsp;&nbsp;go&nbsp;func()&nbsp;{
&nbsp; &nbsp;defer&nbsp;p.wg.Done()

&nbsp; &nbsp;for&nbsp;{
&nbsp; &nbsp;&nbsp;select&nbsp;{
&nbsp; &nbsp;&nbsp;case&nbsp;<-ctx.Done():
&nbsp; &nbsp; &nbsp;log.Printf("worker exit, worker=%d, reason=%v", workerID, ctx.Err())
&nbsp; &nbsp; &nbsp;return

&nbsp; &nbsp;&nbsp;case&nbsp;job, ok := <-p.jobs:
&nbsp; &nbsp; &nbsp;if&nbsp;!ok {
&nbsp; &nbsp; &nbsp; log.Printf("worker exit, worker=%d, reason=job_channel_closed", workerID)
&nbsp; &nbsp; &nbsp;&nbsp;return
&nbsp; &nbsp; &nbsp;}

&nbsp; &nbsp; &nbsp;p.handle(workerID, job)
&nbsp; &nbsp; }
&nbsp; &nbsp;}
&nbsp; }()
&nbsp;}
}

func&nbsp;(p *WorkerPool)&nbsp;Submit(ctx context.Context, job Job)&nbsp;error&nbsp;{
&nbsp;select&nbsp;{
&nbsp;case&nbsp;<-ctx.Done():
&nbsp;&nbsp;return&nbsp;ctx.Err()

&nbsp;case&nbsp;p.jobs <- job:
&nbsp;&nbsp;return&nbsp;nil

&nbsp;default:
&nbsp;&nbsp;return&nbsp;fmt.Errorf("worker pool queue full, job=%s", job.ID)
&nbsp;}
}

func&nbsp;(p *WorkerPool)&nbsp;Stop()&nbsp;{
&nbsp;close(p.jobs)
&nbsp;p.wg.Wait()
}

func&nbsp;(p *WorkerPool)&nbsp;handle(workerID&nbsp;int, job Job)&nbsp;{
&nbsp;start := time.Now()

&nbsp;defer&nbsp;func()&nbsp;{
&nbsp;&nbsp;if&nbsp;r :=&nbsp;recover(); r !=&nbsp;nil&nbsp;{
&nbsp; &nbsp;log.Printf("job panic, worker=%d, job=%s, err=%v", workerID, job.ID, r)
&nbsp; }
&nbsp;}()

&nbsp;// 这里换成真实业务,比如同步订单、清洗文件、调用三方接口
&nbsp;if&nbsp;err := pushToRemote(job); err !=&nbsp;nil&nbsp;{
&nbsp; log.Printf("job failed, worker=%d, job=%s, cost=%s, err=%v",
&nbsp; &nbsp;workerID, job.ID, time.Since(start), err)
&nbsp;&nbsp;return
&nbsp;}

&nbsp;log.Printf("job done, worker=%d, job=%s, cost=%s",
&nbsp; workerID, job.ID, time.Since(start))
}

func&nbsp;pushToRemote(job Job)&nbsp;error&nbsp;{
&nbsp;time.Sleep(80&nbsp;* time.Millisecond)
&nbsp;return&nbsp;nil
}

这个池子没什么花活。

size 控制同时跑多少个 goroutine。

jobs 是任务队列。

Start 只启动固定数量的 worker。

Submit 负责投递任务。

Stop 负责收口。

看着简单,但几个点不能乱动。

第一个点,Submit 里我用了 default

select&nbsp;{
case&nbsp;p.jobs <- job:
&nbsp;return&nbsp;nil
default:
&nbsp;return&nbsp;fmt.Errorf("worker pool queue full, job=%s", job.ID)
}

这个写法的意思是:队列满了就立刻返回错误,不在这里傻等。

这地方我比较保守。因为很多线上卡死,不是 worker 不干活,而是提交任务的 goroutine 全堵在写 channel 上,最后调用链一层一层堵住。

当然,不是所有业务都要立刻失败。

如果是离线导入,可以允许等一会儿:

func&nbsp;(p *WorkerPool)&nbsp;SubmitWait(ctx context.Context, job Job)&nbsp;error&nbsp;{
&nbsp;timer := time.NewTimer(3&nbsp;* time.Second)
&nbsp;defer&nbsp;timer.Stop()

&nbsp;select&nbsp;{
&nbsp;case&nbsp;p.jobs <- job:
&nbsp;&nbsp;return&nbsp;nil

&nbsp;case&nbsp;<-timer.C:
&nbsp;&nbsp;return&nbsp;fmt.Errorf("submit timeout, job=%s", job.ID)

&nbsp;case&nbsp;<-ctx.Done():
&nbsp;&nbsp;return&nbsp;ctx.Err()
&nbsp;}
}

这里我宁愿给一个 3 秒超时,也不愿意无限等。

无限等待这种东西,在代码里看着温柔,在线上经常就是事故入口。

第二个点,worker 里面一定要 recover

不是说业务代码就应该 panic,而是批量任务里只要有一条脏数据把 goroutine 打死,这个 worker 就少一个。少一个还不明显,少到最后池子还在,活没人干,日志也不一定好看。

所以这里要兜一下:

defer&nbsp;func()&nbsp;{
&nbsp;if&nbsp;r :=&nbsp;recover(); r !=&nbsp;nil&nbsp;{
&nbsp; log.Printf("job panic, worker=%d, job=%s, err=%v", workerID, job.ID, r)
&nbsp;}
}()

我不喜欢把 panic 吞得悄无声息。至少要把 workerIDjobID 打出来,不然排查的时候只能靠猜。

第三个点,关闭顺序别写反。

Stop 里先 close(p.jobs),再 Wait

func&nbsp;(p *WorkerPool)&nbsp;Stop()&nbsp;{
&nbsp;close(p.jobs)
&nbsp;p.wg.Wait()
}

不能先等再关。你不关 channel,worker 就一直等任务,Wait 永远不回来。

也不要在多个地方 close 同一个 channel。这个错误很低级,但线上确实见过。一般我会把关闭动作只放在池子自己内部,外面别碰 jobs

用的时候大概这样:

func&nbsp;RunImport(ctx context.Context, rows []string)&nbsp;error&nbsp;{
&nbsp;p := NewWorkerPool(8,&nbsp;64)
&nbsp;p.Start(ctx)
&nbsp;defer&nbsp;p.Stop()

&nbsp;for&nbsp;i, row :=&nbsp;range&nbsp;rows {
&nbsp; job := Job{
&nbsp; &nbsp;ID: &nbsp; &nbsp; &nbsp;fmt.Sprintf("row-%d", i+1),
&nbsp; &nbsp;Payload: row,
&nbsp; }

&nbsp;&nbsp;if&nbsp;err := p.SubmitWait(ctx, job); err !=&nbsp;nil&nbsp;{
&nbsp; &nbsp;log.Printf("submit failed, job=%s, err=%v", job.ID, err)
&nbsp; &nbsp;return&nbsp;err
&nbsp; }
&nbsp;}

&nbsp;return&nbsp;nil
}

这里的 8 不是拍脑袋越大越好。

如果任务主要是 CPU 计算,比如压缩、加密、图片处理,worker 数接近 CPU 核数就差不多了。

如果任务主要是 IO,比如调接口、写数据库、读文件,可以适当放大。但也要看下游扛不扛得住。你本地开 100 个 worker 很爽,数据库连接池只有 20 个,最后还是在连接池那里排队。

这也是我不太喜欢“无脑调大并发”的原因。

并发高,不等于吞吐高。很多时候只是把等待从 A 点挪到 B 点。

一个能用的 goroutine 池,至少要回答四个问题:

任务最多堆多少?

worker 最多跑多少?

提交失败怎么办?

服务退出时,剩下的任务怎么收?

这四个问题没想清楚,就别急着封装成公共库。

Go 的并发工具很锋利,go func() 一行就能起飞,也一行就能把系统打穿。简单资源池的价值不在于代码多高级,而是把“随便起 goroutine”这件事收住。

线上系统很多时候不怕慢一点,怕的是没边界。池子就是这个边界。


免责声明:

本文所载程序、技术方法仅面向合法合规的安全研究与教学场景,旨在提升网络安全防护能力,具有明确的技术研究属性。

任何单位或个人未经授权,将本文内容用于攻击、破坏等非法用途的,由此引发的全部法律责任、民事赔偿及连带责任,均由行为人独立承担,本站不承担任何连带责任。

本站内容均为技术交流与知识分享目的发布,若存在版权侵权或其他异议,请通过邮件联系处理,具体联系方式可点击页面上方的联系我

本文转载自:Go语言教程 go go《分析一个简单的goroutine资源池》

评论:0   参与:  0