直接上生产易卡顿:重试队列未分片致brpoplpush阻塞,且默认不校验payload json合法性,解析panic无原始日志;需手动配置onfailure日志、显式重试策略、payload预校验,并隔离worker消费以保障可靠性。

为什么不用 github.com/hibiken/asynq 直接上生产?
它确实开箱即用,但上线后容易在高并发场景下卡住任务重试逻辑——重试队列没做分片,所有失败任务挤在同一个 Redis List 里,BRPOPLPUSH 阻塞时长会随积压量线性增长。更关键的是,默认不校验 task.Payload 的 JSON 结构合法性,一旦上游传入非法字段(比如 null 嵌套在 required 字段里),worker 解析直接 panic,且错误日志不带原始 payload,排查得翻 Redis raw data。
实操建议:
- 必须覆盖
asynq.ServerOption中的OnFailure回调,手动把task.Type、task.Payload和err.Error()写进结构化日志(比如 Loki + Promtail) - 禁用默认重试策略,在
asynq.Task{}构造时显式传入asynq.MaxRetry(3)和asynq.RetryDelay(5 * time.Second),避免指数退避导致延迟不可控 - 对 payload 做预校验:在
asynq.HandlerFunc入口加一层json.Unmarshal+validator校验,校验失败立刻 return error,不进重试队列
如何让多个 worker 实例真正隔离消费同一队列?
很多人以为只要启动多个 asynq.NewServer 实例就自动负载均衡了,其实不然。Redis 的 List 消费本质是“争抢”,当多个 worker 同时 BRPOPLPUSH 同一个 key,Redis 不保证公平性,会出现某台机器持续拿到任务、另一台长期空转的情况,尤其在任务处理时间差异大时(比如有的任务调第三方 API 耗时 200ms,有的本地计算要 2s)。
实操建议:
- 用
asynq.RedisClientOpt{Addr: "redis://host:port/1"}给每个 worker 分配独立 DB(如 db=1, db=2),再配合asynq.QueueName("queue-a"),让不同实例绑定不同物理队列,而非共享一个 list key - 如果必须共用队列,启用
asynq.ServerOption{Concurrency: 1}+ 自研简单轮询调度器:用 Redis Hash 存每个 worker 的 last_seen 时间戳,主调度器每 5 秒扫一次,把新任务按 last_seen 排序后推给最久未活跃的 worker 对应的专属 list - 禁用
asynq.DefaultServeMux,自己 new 一个 mux 并注册 handler,避免不同业务 task type 混在同一个 handler 里互相影响超时和 panic 恢复
asynq 的 context.Context 超时传递为什么经常失效?
根本原因是 asynq 内部用 context.WithTimeout 创建子 context 时,timeout 是从任务入队时间算起,不是从 worker 开始执行时算起。如果任务在队列里积压了 3 分钟,而你设了 context.WithTimeout(ctx, 30*time.Second),那 worker 真正拿到 context 时剩余时间可能只剩 100ms,导致刚解码完 payload 就被 cancel。
实操建议:
- 永远不要依赖入队时设置的全局 timeout;改用
asynq.Task.GetCustom("deadline_unix_ms")在入队前写入绝对截止时间戳(比如time.Now().Add(30 * time.Second).UnixMilli()),worker 执行时再换算成剩余时间 - 在 handler 开头加
if time.Now().UnixMilli() > deadlineMs { return errors.New("task expired") },比靠 context.Done() 更可靠 - 对需要长时间运行的任务(如导出大文件),改用 “分段提交” 模式:先发一个轻量 task 启动流程并生成 job_id,再由该 task 主动发多个子 task 处理分片,每个子 task 单独设 timeout
Redis 故障时任务不丢的关键动作
默认配置下,asynq 把任务存 Redis List,失败重试也只靠 Redis,一旦 Redis 宕机超过 asynq.ServerOption{RetryDelay: ...} 设置的间隔,正在重试的任务就彻底丢失。这不是理论风险——K8s 集群滚动更新 Redis Pod 时,client 连接闪断常见。
实操建议:
- 开启
asynq.RedisClientOpt{DialTimeout: 3 * time.Second, ReadTimeout: 5 * time.Second, WriteTimeout: 5 * time.Second},避免单点卡死拖垮整个 worker - 在
OnFailure回调里,把失败任务 dump 到本地磁盘(比如/var/log/asynq-fallback/下按天分目录的 JSON 文件),另起一个守护 goroutine 定期扫描该目录,尝试重新client.Enqueue,成功后删文件 - 不要依赖 Redis 持久化(RDB/AOF)保任务——RDB 是定时快照,AOF rewrite 期间可能丢数据;fallback 到本地磁盘才是可控兜底
真正的难点不在怎么选框架,而在怎么让每个失败分支都有可验证的 fallback 路径。比如重试失败后写磁盘,得确保磁盘满时有告警而不是静默失败;比如 deadline 计算,得考虑 worker 所在节点时钟漂移。这些细节不写进监控指标,上线后就是黑盒。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











