如何确保Kafka拒绝重复消息?(ACK丢失重试场景)
这绝对是Kafka生产环境里高频碰到的坑——ACK丢了但消息已经落地,重试就会导致重复。要做到Broker不存重复、消费者看不到重复,得从生产端的幂等性保障入手,再配合消费端的兜底处理,双管齐下:
一、生产端:利用Kafka内置特性+全局唯一ID杜绝重复写入
1. 开启生产者幂等性(最基础的保障)
Kafka从0.11版本开始支持幂等生产者,这是解决单生产者单会话重复的最优解。开启后,生产者会自动为每个消息分配:
- 一个全局唯一的
Producer ID (PID)(生产者启动时生成,重启会变) - 每个分区内递增的
Sequence Number
Broker收到消息时,会记录每个PID对应分区的最大序列号。如果收到的消息PID、分区、序列号都和已写入的重复,就直接丢弃,不会再次写入。
配置方式很简单,在生产者参数里加:
enable.idempotence=true acks=all # 幂等性要求必须等待所有副本确认,否则无法保证一致性 retries=3 # 根据业务设置合理的重试次数
⚠️ 注意:幂等性只覆盖单生产者、单会话、单分区的场景。如果生产者重启(PID变化)、或者消息跨分区发送,这个机制就失效了,这时候就得靠你预先生成的UUID消息ID了。
2. 基于UUID消息ID的全局去重(跨场景兜底)
既然你已经有了预先生成的UUID作为唯一消息ID,可以把这个ID作为消息的Key,或者放在消息的Header里,然后通过两种方式实现全局去重:
方式A:结合外部存储做发送前校验
在生产者发送消息前,先把UUID写入一个分布式存储(比如Redis),用SETNX命令(只有不存在时才写入),同时设置过期时间(和Kafka Topic的消息保留期一致)。如果SETNX返回成功,就发送消息;如果返回失败,说明这个消息已经发送过,直接跳过重试。
这种方式要注意原子性:可以先标记UUID为“待发送”,发送成功后再更新状态为“已发送”,避免因网络波动导致的误判。
方式B:用Kafka事务实现Exactly-Once语义
如果你的业务要求严格的Exactly-Once(比如跨多个分区发送消息,或者需要和消费-生产链路结合),可以开启事务生产者。事务是幂等性的超集,它会把多个消息发送操作包裹在一个事务里,Broker要么全部提交要么全部回滚,同时还能保证跨会话、跨分区的幂等。
配置需要加:
enable.idempotence=true transactional.id=your-unique-transaction-id # 每个生产者实例的事务ID要唯一 acks=all retries=3
二、消费端:兜底处理,确保即使漏网也不会重复消费
哪怕生产端做了100%的保障,也可能因为极端情况(比如Broker故障恢复后的数据不一致)出现重复消息,所以消费端必须自己做幂等处理:
- 核心逻辑:用你预先生成的UUID作为去重键,记录已经成功处理的消息ID。
- 存储选择:
- 本地缓存:比如Guava Cache,适合单实例消费者,设置合理的过期时间(比消息保留期短就行)。
- 分布式存储:比如Redis、MySQL,适合集群消费者,保证所有实例都能共享去重状态。
- 原子性保障:处理业务逻辑和记录消息ID必须在同一个事务里。比如用MySQL的话,把“更新业务数据”和“插入消息ID到去重表”放在同一个事务中,要么都成功,要么都失败,避免出现“业务处理了但没记录ID”导致的重复处理。
关键提醒
- 绝对不要用时间戳作为去重依据!你已经提到重试时时间戳会变,这个完全不可靠。
- 优先用Kafka内置的幂等性和事务,这是最省心的方案,自己写的去重逻辑容易引入性能瓶颈或者分布式一致性问题。
内容的提问来源于stack exchange,提问作者Drathier

