redis.setnx易漏判重复消息因缺乏原子性:多个实例同时查key不存在后均执行set,导致重复处理;正确做法是用带ttl的原子setnx或lua脚本实现“查+写”一体。

为什么用 Redis.SetNX 容易漏判重复消息
因为 SetNX 单独调用不带过期时间,或拆成 Exists + Set 两步,中间存在竞态窗口:多个服务实例同时查到 key 不存在,全都会写入并执行业务逻辑。线上最常见错误是只写 rdb.SetNX(ctx, key, "1", 0),TTL 设为 0 就等于没过期,机器宕机后 key 永久卡住。
正确做法必须用原子命令:rdb.SetNX(ctx, key, traceID, 300*time.Second)(300 秒 ≥ 消息处理最长耗时 + 余量),且 value 填 traceID 而非空字符串,方便日志对齐和人工排查。
- key 格式推荐:
"idempotent:msg:{service_name}:{topic}:{msg_id_hash}",避免不同服务/主题 key 冲突 - 不要用客户端传的 raw
msg_id直接当 key——可能含非法字符或超长,先做sha256.Sum256哈希再截取前 16 字节 - 如果 Redis 不可用,
err != nil时不能静默放行,应返回503 Service Unavailable或走降级队列,否则幂等语义彻底失效
怎么用 Lua 脚本保证“查+写”绝对原子
Go 的 Redis 客户端(如 github.com/go-redis/redis/v9)不直接暴露 EVAL,但你可以封装一个 rdb.Eval(ctx, script, []string{key}, ttlSeconds) 调用。脚本内容必须是单次原子执行:
if redis.call("GET", KEYS[1]) == false then
redis.call("SET", KEYS[1], "1", "EX", ARGV[1])
return 1
else
return 0
end
这个脚本在 Redis 单线程内完成判断与写入,彻底堵死并发漏判。注意:ARGV[1] 是 TTL 秒数,不是 time.Duration;KEYS[1] 是完整 key 名,别拼错。
- 别在 Lua 里做复杂逻辑(比如解析 JSON 或调外部 API),它只负责“是否已处理”的判定
- 脚本返回 1 表示首次处理,0 表示已存在——业务层据此跳过后续逻辑,而不是靠 error 判断
- 如果消息体含重试标识(如
retry_count),可在脚本里读取旧值并递增,但需额外用GETSET配合,增加复杂度,一般不建议
数据库唯一索引为什么不能只建在 msg_id 上
只对 msg_id 建 UNIQUE KEY 看似简单,但会漏掉同一消息被投递到不同业务表的场景。比如支付消息既写 payments 表,又触发发券写 vouchers 表——两个表各自唯一,但发券动作仍可能被执行两次。
真正兜底的索引必须覆盖业务语义:UNIQUE KEY uk_service_topic_msg (`service`, `topic`, `msg_id_hash`)。其中 service 是消费方服务名(如 "order-service"),topic 是 Kafka 主题或 RocketMQ 的 tag,msg_id_hash 是哈希后的消息 ID。
- 插入这条幂等记录的动作,必须和主业务逻辑(如创建订单)在同一个 DB 事务里提交,否则无法保证原子性
- 捕获冲突错误要精准:
if pgx.ErrCodeUniqueViolation(err) || mysql.IsDupEntryError(err),别用strings.Contains(err.Error(), "duplicate")这种脆弱匹配 - 不要在事务里做 HTTP 调用或 sleep——锁持有时间越长,DB 并发越差
context.WithValue 透传幂等信息容易踩什么坑
很多人把 traceID 或 msg_id 存进 context.WithValue(ctx, idempotentKey, val),然后在 handler 里取。问题在于:如果消息经过多级服务中转(A→B→C),而 B 没有显式把 context 透传给 C 的 RPC 调用,那 C 就拿不到原始幂等键。
更稳妥的方式是:从消息头(如 Kafka 的 headers["X-Trace-ID"] 或 gRPC 的 metadata)里提取,并在每一跳都显式注入到下游 context 中。Gin/Echo 中间件可以统一做这事,但 net/http 的 HandlerFunc 必须手动传递。
- 别用闭包变量存幂等 key——goroutine 并发下会互相污染
- value 类型必须是导出的自定义类型(如
type IdempotentKey string),不能用string作 key,否则不同包之间可能冲突 - 如果消息本身不带 traceID,就由第一个接入服务生成并写回消息头,后续所有跳都复用它,而不是每跳都新生成











