文章总结: 针对Go无节制创建goroutine导致线上卡死的问题,本文指出需建立并发边界。作者基于channel实现轻量级资源池,强调三处关键设计:任务提交须快速失败或设超时防阻塞,worker内须recover防脏数据打挂,关闭时须先close再Wait。建议按任务特性配置worker数,通过池化约束并发以保障系统稳定。 综合评分: 91 文章分类: 其他
分析一个简单的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 {
if size <= 0 {
size = 1
}
if queueSize <= 0 {
queueSize = size * 2
}
return &WorkerPool{
size: size,
jobs: make(chan Job, queueSize),
}
}
func (p *WorkerPool) Start(ctx context.Context) {
for i := 0; i < p.size; i++ {
workerID := i + 1
p.wg.Add(1)
go func() {
defer p.wg.Done()
for {
select {
case <-ctx.Done():
log.Printf("worker exit, worker=%d, reason=%v", workerID, ctx.Err())
return
case job, ok := <-p.jobs:
if !ok {
log.Printf("worker exit, worker=%d, reason=job_channel_closed", workerID)
return
}
p.handle(workerID, job)
}
}
}()
}
}
func (p *WorkerPool) Submit(ctx context.Context, job Job) error {
select {
case <-ctx.Done():
return ctx.Err()
case p.jobs <- job:
return nil
default:
return fmt.Errorf("worker pool queue full, job=%s", job.ID)
}
}
func (p *WorkerPool) Stop() {
close(p.jobs)
p.wg.Wait()
}
func (p *WorkerPool) handle(workerID int, job Job) {
start := time.Now()
defer func() {
if r := recover(); r != nil {
log.Printf("job panic, worker=%d, job=%s, err=%v", workerID, job.ID, r)
}
}()
// 这里换成真实业务,比如同步订单、清洗文件、调用三方接口
if err := pushToRemote(job); err != nil {
log.Printf("job failed, worker=%d, job=%s, cost=%s, err=%v",
workerID, job.ID, time.Since(start), err)
return
}
log.Printf("job done, worker=%d, job=%s, cost=%s",
workerID, job.ID, time.Since(start))
}
func pushToRemote(job Job) error {
time.Sleep(80 * time.Millisecond)
return nil
}
这个池子没什么花活。
size 控制同时跑多少个 goroutine。
jobs 是任务队列。
Start 只启动固定数量的 worker。
Submit 负责投递任务。
Stop 负责收口。
看着简单,但几个点不能乱动。
第一个点,Submit 里我用了 default。
select {
case p.jobs <- job:
return nil
default:
return fmt.Errorf("worker pool queue full, job=%s", job.ID)
}
这个写法的意思是:队列满了就立刻返回错误,不在这里傻等。
这地方我比较保守。因为很多线上卡死,不是 worker 不干活,而是提交任务的 goroutine 全堵在写 channel 上,最后调用链一层一层堵住。
当然,不是所有业务都要立刻失败。
如果是离线导入,可以允许等一会儿:
func (p *WorkerPool) SubmitWait(ctx context.Context, job Job) error {
timer := time.NewTimer(3 * time.Second)
defer timer.Stop()
select {
case p.jobs <- job:
return nil
case <-timer.C:
return fmt.Errorf("submit timeout, job=%s", job.ID)
case <-ctx.Done():
return ctx.Err()
}
}
这里我宁愿给一个 3 秒超时,也不愿意无限等。
无限等待这种东西,在代码里看着温柔,在线上经常就是事故入口。
第二个点,worker 里面一定要 recover。
不是说业务代码就应该 panic,而是批量任务里只要有一条脏数据把 goroutine 打死,这个 worker 就少一个。少一个还不明显,少到最后池子还在,活没人干,日志也不一定好看。
所以这里要兜一下:
defer func() {
if r := recover(); r != nil {
log.Printf("job panic, worker=%d, job=%s, err=%v", workerID, job.ID, r)
}
}()
我不喜欢把 panic 吞得悄无声息。至少要把 workerID、jobID 打出来,不然排查的时候只能靠猜。
第三个点,关闭顺序别写反。
Stop 里先 close(p.jobs),再 Wait。
func (p *WorkerPool) Stop() {
close(p.jobs)
p.wg.Wait()
}
不能先等再关。你不关 channel,worker 就一直等任务,Wait 永远不回来。
也不要在多个地方 close 同一个 channel。这个错误很低级,但线上确实见过。一般我会把关闭动作只放在池子自己内部,外面别碰 jobs。
用的时候大概这样:
func RunImport(ctx context.Context, rows []string) error {
p := NewWorkerPool(8, 64)
p.Start(ctx)
defer p.Stop()
for i, row := range rows {
job := Job{
ID: fmt.Sprintf("row-%d", i+1),
Payload: row,
}
if err := p.SubmitWait(ctx, job); err != nil {
log.Printf("submit failed, job=%s, err=%v", job.ID, err)
return err
}
}
return nil
}
这里的 8 不是拍脑袋越大越好。
如果任务主要是 CPU 计算,比如压缩、加密、图片处理,worker 数接近 CPU 核数就差不多了。
如果任务主要是 IO,比如调接口、写数据库、读文件,可以适当放大。但也要看下游扛不扛得住。你本地开 100 个 worker 很爽,数据库连接池只有 20 个,最后还是在连接池那里排队。
这也是我不太喜欢“无脑调大并发”的原因。
并发高,不等于吞吐高。很多时候只是把等待从 A 点挪到 B 点。
一个能用的 goroutine 池,至少要回答四个问题:
任务最多堆多少?
worker 最多跑多少?
提交失败怎么办?
服务退出时,剩下的任务怎么收?
这四个问题没想清楚,就别急着封装成公共库。
Go 的并发工具很锋利,go func() 一行就能起飞,也一行就能把系统打穿。简单资源池的价值不在于代码多高级,而是把“随便起 goroutine”这件事收住。
线上系统很多时候不怕慢一点,怕的是没边界。池子就是这个边界。
免责声明:
本文所载程序、技术方法仅面向合法合规的安全研究与教学场景,旨在提升网络安全防护能力,具有明确的技术研究属性。
任何单位或个人未经授权,将本文内容用于攻击、破坏等非法用途的,由此引发的全部法律责任、民事赔偿及连带责任,均由行为人独立承担,本站不承担任何连带责任。
本站内容均为技术交流与知识分享目的发布,若存在版权侵权或其他异议,请通过邮件联系处理,具体联系方式可点击页面上方的联系我。
本文转载自:Go语言教程 go go《分析一个简单的goroutine资源池》
版权声明
本站仅做备份收录,仅供研究与教学参考之用。
读者将信息用于其他用途的,全部法律及连带责任由读者自行承担,本站不承担任何责任。










评论