主题
24 · Asynq 任务队列
目标:用 Asynq 在 Redis 上实现「入队 → Worker 消费 → 重试 / 延时」,对标 Node 的 Bull。
前置:20 · go-redis · 23 · RabbitMQ(概念对照)
官方:hibiken/asynq · 文档 Wiki
1. Asynq vs RabbitMQ:怎么选
| 维度 | Asynq | RabbitMQ |
|---|---|---|
| 依赖 | 已有 Redis | 独立 Broker |
| 心智 | Job:类型 + Payload + 重试 | AMQP:Exchange/Queue |
| 适合 | 同进程族异步任务(邮件、导出、压缩) | 多语言、复杂路由、强解耦 |
| Node 对照 | Bull / BullMQ | amqplib |
练手仓:同一仓库里的异步活优先 Asynq;跨服务事件总线再 Rabbit/Kafka。
2. 概念速览
text
Client.Enqueue(Task) → Redis → Server + Handler(mux)
↓
成功 / 失败重试 / 死信(archived)1
2
3
2
3
| 能力 | 用途 |
|---|---|
ProcessIn / ProcessAt | 延时任务 |
MaxRetry | 失败重试次数 |
Queue("critical") | 多队列优先级 |
| Scheduler | 周期入队(可与 25 Cron 分工) |
| asynqmon | Web 看任务状态(可选) |
3. 最小可跑 demo
bash
go get github.com/hibiken/asynq1
任务定义:
go
const TypeExportPDF = "note:export_pdf"
type ExportPayload struct {
NoteID string `json:"noteId"`
}
func NewExportTask(noteID string) (*asynq.Task, error) {
b, err := json.Marshal(ExportPayload{NoteID: noteID})
if err != nil {
return nil, err
}
return asynq.NewTask(TypeExportPDF, b), nil
}1
2
3
4
5
6
7
8
9
10
11
12
13
2
3
4
5
6
7
8
9
10
11
12
13
HTTP / 业务侧入队:
go
client := asynq.NewClient(asynq.RedisClientOpt{Addr: "127.0.0.1:6379"})
defer client.Close()
task, _ := NewExportTask("n1")
info, err := client.Enqueue(task,
asynq.MaxRetry(3),
asynq.Timeout(2*time.Minute),
asynq.Queue("default"),
)
_ = info1
2
3
4
5
6
7
8
9
10
2
3
4
5
6
7
8
9
10
Worker:
go
srv := asynq.NewServer(
asynq.RedisClientOpt{Addr: "127.0.0.1:6379"},
asynq.Config{Concurrency: 4, Queues: map[string]int{"default": 1}},
)
mux := asynq.NewServeMux()
mux.HandleFunc(TypeExportPDF, func(ctx context.Context, t *asynq.Task) error {
var p ExportPayload
if err := json.Unmarshal(t.Payload(), &p); err != nil {
return err // 返回 error → 按策略重试
}
log.Printf("export %s", p.NoteID)
return nil
})
if err := srv.Run(mux); err != nil {
log.Fatal(err)
}1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
本地:Compose 已有 Redis 即可;API 进程 Enqueue,另起一个 cmd/worker。
4. 延时与幂等
go
_, _ = client.Enqueue(task, asynq.ProcessIn(10*time.Minute))1
幂等:同一 noteId 重复入队可能跑两次 → Handler 内用 Redis SetNX / DB 唯一状态机(pending|done)。
5. 动手清单
- [ ] Client 入队 + Worker 打印 payload
- [ ] Handler
return errors.New("boom"):看到重试 - [ ]
ProcessIn(5*time.Second):到期才执行 - [ ] 两个队列
critical/default,权重不同 - [ ] 能讲清:与 Rabbit 何时二选一
6. 项目驱动
| 场景 | 做法 |
|---|---|
| 导出 PDF | 接口 202 + task_id;Worker 写 MinIO(26) |
| 图片压缩 | 上传回调后入队 |
| AI 总结 | 长任务 + Timeout + 进度写 Redis |
| 发邮件 | 失败自动重试,避免拖垮请求线程 |
7. 常见坑 + AI 审查
| 坑 | 说明 |
|---|---|
| Handler 吞掉 error | 不 return err → 永不重试 |
| 与业务共用一个 Redis DB 且 key 冲突 | 可用不同 DB index 或 key 前缀 |
| Worker 不优雅退出 | 关注 Shutdown / 信号处理 |
| AI 把「定时」全塞进 Asynq Scheduler | 简单 cron 可用 25;复杂再 Scheduler |
| Payload 塞超大文件 | 只传 ID;内容放对象存储 |
8. 下一篇
进程内 / 分布式定时触发:
→ 25 · Cron 定时任务
