beego的beego/queue需显式导入并调用queue.start()启用,否则任务堆积不执行;默认内存队列重启丢失,生产环境应配置redis(addr须带redis://前缀);任务函数签名必须为func() error,参数须json序列化为[]byte传入;需手动设置queue.setmaxconcurrent()控制并发数,且务必在start()前调用。

Beego内置beego/queue怎么启用和配置
Beego的beego/queue不是开箱即用的模块,必须显式导入并初始化。它本身不启动后台协程,需要你手动调用Start(),否则任务会一直堆积在内存队列里不动。
常见错误是只注册了任务但没调用queue.Start(),结果发了100个任务,日志里完全没反应——其实全卡在未消费状态。
- 必须在
main()或init()中调用queue.Start(),且只能调用一次 - 默认使用内存队列(
memory),重启后任务丢失;生产环境务必切换为redis或rabbitmq - 配置Redis队列时,
addr参数必须带redis://前缀,例如"redis://127.0.0.1:6379/0",漏掉协议会导致连接失败但无明确报错 - 任务函数签名必须是
func() error,返回nil才视为成功;若panic或返回非nil error,任务会被重试(默认3次)
queue.Push()传参要注意什么
Beego队列不支持直接传递结构体、指针或闭包,所有参数必须序列化为[]byte。你不能写queue.Push(myHandler, user),因为user会被强制转成fmt.Sprintf("%v", user),容易丢失字段或引发panic。
正确做法是提前把参数JSON序列化,再传入:
data, _ := json.Marshal(map[string]interface{}{"uid": 123, "action": "send_email"})
queue.Push("send_email_task", data)
- 任务名(第一个参数)是字符串标识,用于路由到对应处理器,不要含空格或特殊符号
- 第二个参数必须是
[]byte,别传string——Go里string和[]byte底层不同,直接传会触发隐式转换但可能截断二进制数据 - 如果参数里有时间类型(如
time.Time),JSON序列化后是字符串,反序列化时需手动转回,别依赖json.Unmarshal自动还原
如何避免任务重复执行或丢失
内存队列(memory)下,进程崩溃=全部任务丢失;Redis队列下,若消费者在处理中宕机,未ack的任务会重新入队——这是双刃剑:保障不丢,但也可能重复。
- Redis模式下,
queue.Pop()取出任务后会自动加锁(默认30秒),超时未完成则释放锁并重入队。可通过queue.SetLockTimeout(60)延长 - 关键业务必须实现幂等:比如发邮件任务,先查DB确认是否已发,再执行;或用Redis
SETNX打标记 - 不要在任务函数里直接调用
os.Exit()或log.Fatal(),这会让整个worker退出,后续任务全卡住 - 批量任务慎用
queue.PushMany():它一次性推入多个,但底层仍是逐个消费,无法保证原子性
为什么queue.AddWorker()设了5个却只跑2个协程
Beego队列的worker数量不是并发上限,而是“最多同时拉取多少个待处理任务”。实际并发数还受queue.SetMaxConcurrent()控制,默认值是2。
现象:你调了queue.AddWorker(5),但监控显示CPU只跑2核,日志里同一秒只打印2条“task done”——大概率是没调queue.SetMaxConcurrent(5)。
-
AddWorker(n)只是启动n个goroutine去轮询队列,但每个goroutine默认每次只取1个任务,取完就等下一轮 -
SetMaxConcurrent(n)才是控制“最多几个任务并行执行”,必须在Start()前设置才生效 - 高IO任务(如HTTP请求)建议设高些(如10–20);CPU密集型任务建议等于CPU核心数,避免调度开销
- 注意Goroutine泄漏:如果任务函数里启了子goroutine但没做
sync.WaitGroup或context控制,worker退出后它们还在后台跑
SetMaxConcurrent()和Start()的调用时机——这两处写错,整个队列就形同虚设。











