kafka消费者天然不幂等,因其at-least-once语义使重复消费成为常态,而enable.idempotence仅作用于生产者端、单会话单分区,对消费者无感知;可靠方案是用redis setnx原子校验消息id,要求key含可复现业务上下文、ex时间≥业务最大耗时×2.5、失败时立即ack并return。

为什么 Kafka 消费者天然不幂等
Kafka 的 at-least-once 语义决定了重复消费不是异常,而是常态。消费者崩溃、rebalance、网络超时、手动重置 offset —— 这些都会导致同一条消息被多次拉取。而 enable.idempotence=true 只作用于生产者端,且仅保障单会话、单分区内的写入不重复;它对消费者完全无感知,也不校验业务逻辑是否已执行。
用 Redis SETNX 做前置校验才是可靠方案
核心是原子性:不能先 GET 再 SET,必须用 SET key value EX seconds NX 一步完成“检查 + 占坑”。推荐用 github.com/go-redis/redis/v9 的 rdb.SetNX(ctx, key, "done", ttl)。
-
key必须含业务上下文,例如:"idempotent:order_created:789:1a2b3c",其中789是order_id,1a2b3c是sha256(msg.Body)[:3]的十六进制前缀;纯message_id或拼接routing_key都不可靠 -
EX时间要 ≥ 单次业务最大耗时 × 2.5,比如处理最长 12s,设40 * time.Second -
SetNX返回false时,必须立刻msg.Ack()并return,不能panic,也不能msg.Nack(false, true),否则会触发框架重试
消息 ID 构造必须可复现且防碰撞
生产者和消费者必须用同一套规则生成 ID,否则去重完全失效。常见错误包括:
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
- 依赖
msg.Headers["X-Trace-ID"]或time.Now().UnixNano()—— 每次哈希值都不同 - 只用
msg.Value原始字节做 MD5 —— 若 body 含随机字段、时间戳、traceID,就失去可复现性 - AMQP 的
MessageId默认为空或复用,RabbitMQ/Kafka 都不强制设置,不能直接拿来当唯一键
正确做法是提取业务主键(如 order_id)+ 稳定字段(如 event_type),再加盐哈希,例如:sha256(fmt.Sprintf("%s:%s:%s", eventType, orderId, "v2"))[:12]。
别把数据库唯一约束当幂等依据
靠 INSERT INTO orders (id) VALUES (?) ON CONFLICT DO NOTHING 或主键冲突拦截,只是兜底,不是幂等实现。
- DB 报错发生在副作用之后:扣款、发短信、调外部接口可能已完成
- 高并发下唯一索引冲突会引发行锁等待,拖慢吞吐,甚至死锁
- 分库分表场景中,唯一约束无法跨物理库生效
- 某些 ORM 自动生成 UUID 或时间戳作为主键,根本没业务含义,去重失效
真正的幂等必须在业务逻辑执行前完成判断,且失败路径要干净:log 记录、ack 消息、立即 return —— 不留中间态,不触发重试。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










