必须用lua脚本原子执行查删标记三步,因zrangebyscore+zrem分离调用会导致重复消费或任务丢失;zpopmin是redis 5.0+原生解法;score须统一为秒级整数并校准redis服务器时间。

直接用 redis.ZAdd + ZRangeByScore 轮询是最易上手的方案,但不加原子封装必然重复消费或丢任务——必须用 Lua 脚本把「查、删、标记」三步锁死在 Redis 端执行。
为什么 ZRANGEBYSCORE + ZREM 分开调用会翻车
两个命令之间存在时间窗口:Worker A 查到 3 条到期消息,还没来得及 ZRem,Worker B 同样查到这 3 条,结果两边都开始处理。更糟的是,若 A 处理失败又没重入逻辑,B 也没拿到,任务就永久丢失。
-
ZPOPMIN是 Redis 5.0+ 的原生解法,单命令弹出最小 score 元素,天然避免查删分离 - 老版本只能靠 Lua:先
ZRangeByScore查最小值,再ZRem删除它,整个流程在服务端原子执行 - 别信“先删后处理”的简化逻辑——分布式环境下几毫秒偏差就足以让两个 worker 同时拿到同一任务
score 必须统一为秒级整数,且 Redis 服务器时间要校准
ZRangeByScore 返回空但消息该执行了?八成是 score 单位不一致或时区漂移。Go 侧必须用 time.Now().Unix(),Redis 侧也得按秒算,不能混用毫秒(比如 time.Now().UnixMilli() 存、却用秒查)。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
- 检查 Redis 服务器时间:
redis-cli time对比本地date +%s,偏差超 1 秒就得同步 NTP - 所有
ZAdd的score必须是整数秒,别传浮点数(float64(time.Now().Unix())安全,float64(time.Since().Seconds())有精度风险) -
Max参数动态传strconv.FormatInt(time.Now().Unix(), 10),别硬写死"1714000000"
轮询不能写成 for { zrange; sleep(100ms) },得动态休眠 + 超时防护
固定短间隔轮询既浪费连接又压 Redis;空转时休眠太短,有任务时又响应慢。核心是用 ZCount 预判是否有活,再决定是否拉取,并控制休眠节奏。
- 用
redis.Client.ZRangeByScoreWithScores,传入&redis.ZRangeBy{Min: "-inf", Max: fmt.Sprintf("%d", time.Now().Unix())} - 每次轮询前加
context.WithTimeout(ctx, 3*time.Second),防止单次命令卡死 - 首次启动时用
ZCARD看队列长度,若 > 0 则立刻触发一轮拉取,避免首屏延迟 - 连续 5 次没拿到任务,休眠从 100ms 升到 500ms;一旦有任务,立刻回落到 100ms
Lua 脚本必须同时完成“领取 + 标记”,否则 worker 崩溃任务就卡死
只弹出还不够——得立刻标记“已被某 worker 领取”,否则进程崩溃后任务无法被其他 worker 接管。推荐用 Lua 把 ZPOPMIN 和 HSET 绑定为一个原子操作,写入 processing hash 记录 worker ID 和开始时间。
- 脚本示例中
KEYS[2]是"delay_processing",ARGV[1]是当前 worker ID - 消费成功后,用
HDEL delay_processing <payload></payload>清理状态 - 失败则
ZADD delay_queue <original_score><payload></payload></original_score>重入,但需区分错误类型:临时性错误走指数退避(+1s、+2s、+4s…),上限 5 分钟;参数类错误直接丢弃
真正难的不是写对那几行 Lua,而是所有 worker 共享同一套时间戳单位、同一套时区基准、同一套失败分类逻辑——差一点,就是重复投递或永久丢失。










