
Kafka 的重复消费问题在面试里是一个几乎绕不开的考点在生产环境里也是一个经常出现的高频事故。很多开发者对 Kafka 的认知停留在“削峰填谷、异步解耦”但一旦工作负载从单机版本过渡到多消费者、多分区的真实集群就会发现消息被重复消费的情况远比想象中常见。要从根上解决这个问题先要明白 Kafka 为什么允许重复然后才能在消费端设计出一套防重复机制。下面按“根因、场景、方案、排查、落地”这条主线一步步推进讲清楚 Kafka 如何避免重复消费以及实际项目中应该优先做什么。1. 先想清楚Kafka 为什么会重复消费1.1 从一条消息的完整旅程看 offset 的作用在 Kafka 里一条消息从生产到消费大致要经过这几个环节生产者把消息写入某个 Topic 的某个分区分区给消息分配递增的 offset分区副本在 ISR 列表中保持同步消费者通过 poll 从分区拉取消息消费者处理完成后把自己读到的 offset 提交给 Kafka。offset 是理解重复消费的钥匙。每个分区有一组关键位置生产者写入位置、消费者当前拉取位置、当前提交位置。生产者写入消息消息追加到日志末端消费者拉取消息拉取位置向前移动消费者提交 offset提交位置才向前移动。这里容易产生一个误解消费者拉取到了消息并不代表 Kafka 认为你已经消费成功了。Kafka 判定“这条消息不会再发给同一个组里的其他消费者”靠的是提交 offset而不是业务代码执行完毕。也就是说业务处理完成和 offset 提交完成是两个完全独立的事件。1.2 at least onceKafka 对消息传递的默认保证Kafka 官方文档把消息投递语义分成三种at most once、at least once、exactly once。at most once 表示最多投递一次极端情况下消息会丢at least once 表示至少投递一次极端情况下消息会重复exactly once 表示精确一次需要外部系统配合才能实现。绝大多数生产环境使用的默认组合是 at least once只要消息已经被成功写入消费者重启、副本切换、提交失败等场景都可能导致同一条消息被再次投递。这不是 Kafka 的缺陷而是它在“不丢消息”和“不重复消息”之间选择了前者。与其想方设法让 Kafka 保证不重复不如承认一个前提重复消费是常态消费端必须具备重复消费能力。注意Kafka 默认不保证不重复。使用 Kafka 做业务消息时如果业务本身对重复敏感就必须自己处理幂等。1.3 重复消费的真正源头提交窗口把“处理业务”和“提交 offset”放在时间轴上就能看到重复消费的窗口。假设消费者 A 在 t1 时刻拉取到 offset100 的消息在 t2 时刻完成了业务处理在 t3 时刻提交了 offset。那么如果 A 在 t2 之后、t3 之前崩溃offset 还是 99Kafka 会把 100 再次投递给组内其他消费者或重启后的 A如果 A 在 t3 提交成功后崩溃100 已经提交不会重复如果 A 处理完成但提交请求发送失败offset 仍然停留在 99重启后会从 100 重新消费。所以说重复消费不只是“代码写得不对”而是分布式环境下“处理已完成、提交未完成”这个时间窗口造成的必然结果。搞清楚这一点才能明白为什么后面的每一种方案都在缩小这个窗口或者接受这个窗口并让重复变得无害。2. 四个触发重复消费的真实场景2.1 场景一消费者组重平衡时分区重新分配重平衡是 Kafka 消费者组非常典型的行为。只要消费者组成员发生变化例如新消费者加入、消费者宕机、消费者主动退出消费者组就会重新协调分区分配关系。重平衡过程中消费者 A 已经处理完某条消息但没来得及提交 offset分区被分配给了消费者 BB 会从该分区上次提交的 offset 继续拉取于是这条消息被再次消费。线上常见的触发原因包括消费者进程发布重启、实例被容器调度到其他机器、消费者线程执行超时被判定失联。这些操作如果不做优雅下线处理都可能导致分区被重新分配从而放大重复消费的概率。2.2 场景二业务处理完成但 offset 提交失败手动提交 offset 的场景下提交操作本身也可能失败。比如提交请求超时、网络抖动、session 过期等。此时业务已经写库成功但 Kafka 记录到的 offset 还停留在旧位置。下一次 poll 会重新拿到同一批消息。这种情况比重平衡更加隐蔽因为日志里往往没有明显的异常只有消费端能看到业务重复执行。要识别这种场景需要在提交失败的回调或日志中留下足够信息。例如在 commitAsync 的回调中记录 offset、分区和异常原因不能只打一行 “commit failed” 就结束。2.3 场景三消费处理耗时超过心跳间隔Kafka 消费者组依赖心跳机制感知消费者的存活状态。如果一条消息的处理时间超过了 max.poll.interval.ms比如调用外部接口超时、写慢 SQL、处理大文件消费者会在这个间隔内没有 poll 也没有提交broker 会判定该消费者已经死亡把它移出消费者组并触发重平衡。但此时该消费者所在的线程可能还在继续执行业务。重平衡后未提交 offset 的分区由其他消费者接管消息会被再次处理。这个场景在同步调用外部系统时尤其常见。消费者到数据库、Redis、第三方接口的耗时一旦出现长尾就会超过 Kafka 默认的 300 秒上限。如果业务代码还把异常吞掉了消费者既没有提交 offset也没有退出问题会更加隐蔽。2.4 场景四启用了异步提交但失败后没有补偿机制使用 commitAsync() 提交 offset 时如果提交失败客户端默认不会自动重试只会把异常交给回调函数。如果回调里只是打了日志没有做任何补偿offset 就停留在旧位置。下一轮消费会重复。尤其当代码是“先 commitAsync再继续 poll”的结构时提交失败往往会在很久之后才被注意到。异步提交并不是不能使用而是必须明确它的失败处理策略。对于严格场景可以退回到同步提交如果仍要异步也要在回调里记录失败 offset并配合监控报警避免静默失败。这四个场景共同说明一个道理重复消费会发生在“处理完成”和“提交成功”之间的所有失败路径上想要彻底避免只能让重复消费的业务副作用变为零或者缩短失败路径并保证提交可靠性。3. 方案一在数据库侧做幂等设计3.1 为什么幂等是避免重复消费的第一层防线既然重复消费无法在 Kafka 层做到绝对杜绝最稳妥的做法是让消费逻辑具备幂等性。所谓幂等就是同一条消息被处理任意多次最终产生的业务状态只有一份。比如扣款场景消费者收到一条“用户支付 100 元”的消息如果处理了两次用户就会被扣 200 元这是不能接受的。幂等设计需要保证同一个支付单号只能被扣一次。把业务处理的结果和“幂等记录”绑定在一起重复消息在进入业务处理前就会被识别出来。3.2 幂等表方案一个常见的做法是引入消费记录表。表结构可以这样设计CREATE TABLE kafka_consume_record ( id BIGINT AUTO_INCREMENT PRIMARY KEY, message_key VARCHAR(128) NOT NULL COMMENT 业务幂等键例如订单号, topic VARCHAR(64) NOT NULL, partition_id INT NOT NULL, offset_value BIGINT NOT NULL, consume_status TINYINT NOT NULL DEFAULT 0 COMMENT 0处理中, 1处理完成, consume_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_message_key (message_key) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;这里的 message_key 是业务幂等键它不一定等于 Kafka 的 offset。同一业务消息可能由生产者重试多次每次携带不同的 offset但业务唯一标识是稳定的。用业务 key 去重比单纯用 topic partition offset 更可靠。消费逻辑先尝试插入幂等记录插入成功才继续业务处理public void consume(ConsumerRecordString, String record) { String messageKey record.key(); // 如果消息没有业务 key可以从消息体解析 if (messageKey null || messageKey.isEmpty()) { messageKey buildMessageKeyFromBody(record.value()); } int exists consumeRecordMapper.countByMessageKey(messageKey); if (exists 0) { log.info(duplicate consume, key{}, messageKey); return; } ConsumeRecord recordDO new ConsumeRecord(); recordDO.setMessageKey(messageKey); recordDO.setTopic(record.topic()); recordDO.setPartitionId(record.partition()); recordDO.setOffsetValue(record.offset()); try { consumeRecordMapper.insert(recordDO); } catch (DuplicateKeyException e) { log.info(duplicate key conflict, key{}, messageKey); return; } try { doBusiness(record.value()); } catch (Exception e) { log.error(consume business error, e); // 根据业务决定是否重试不能直接吞掉异常 } }这段代码的关键点先通过唯一键插入尝试用数据库的唯一约束兜底竞态条件。简单的“先查再插”在高并发下仍有问题两个消费者同时查不到记录然后同时插入如果没有唯一约束就会产生两条记录。唯一约束保证了只有一条能成功。3.3 幂等表在实践中的取舍数据库幂等方案实现成本低适合已有 MySQL 的团队。但它有一个问题幂等表会不断增长必须要有清理策略。可以按保留时间定期删除已经过期的消费记录也可以按业务维度把“已经完成的消费记录”归档到历史库。另一个容易忽略的问题是先插幂等记录再执行业务如果业务执行失败这条记录只能删除或标记为失败否则后续重试消息会被“幂等”挡在外面业务永远无法补偿。推荐的做法是insert 后把状态标记为“处理中”业务成功后再更新为“处理完成”重试时看到“处理中”的记录可以结合超时时间决定继续等待还是重新执行。注意幂等表和业务表最好放在同一个数据库事务里否则会出现“业务成功写入、幂等记录没写入”的反向不一致导致下一次重复消息又把业务执行一遍。4. 方案二用 Redis 做轻量级去重4.1 SETNX 命令的语义Redis 的 setnx 是 set if not exists 的缩写。执行 setnx key value 时如果 key 不存在设置成功并返回 1如果 key 已经存在什么都不做并返回 0。这一语义天然适合做消费幂等。消费前先执行 setnx如果返回 1说明当前消息是第一次消费继续执行业务如果返回 0说明已经处理过直接跳过。4.2 Redis 幂等的代码片段public boolean consumeIfAbsent(String messageKey) { String key kafka:consume: messageKey; Boolean first redisTemplate.opsForValue().setIfAbsent(key, 1, Duration.ofHours(2)); if (first null || !first) { log.info(duplicate consume by redis, key{}, messageKey); return false; } return true; }这里需要注意 setIfAbsent 可以带过期时间避免 key 无限占用 Redis 内存。过期时间要根据业务重试窗口设置如果业务允许 7 天内重复消息被丢弃那过期时间至少要覆盖 7 天如果只是防止秒级或分钟级重复两小时通常足够。Redis 方案更适合高吞吐、不依赖事务性幂等表的场景。它的问题是如果 Redis 本身发生故障或者 key 提前过期重复消息仍然会穿透到业务层。所以 Redis 方案通常只作为第一道拦截不能作为唯一防线。4.3 Redis 方案与数据库方案的选择对比维度数据库幂等表Redis 去重实现成本需要建表、写 Mapper 和事务只需一个 setnx 调用一致性可以和业务事务合并与业务操作天然分离性能每一条消息多一次数据库写每一条消息多一次 Redis 写防穿透唯一约束兜底可靠依赖键值和过期时间设置数据清理需要定期清理或归档通过过期时间自动清理适用场景强一致、对精确性要求高高并发、允许短暂重复触发一次拦截根据场景可以组合使用先用数据库唯一约束做最终兜底再在消费入口用 Redis 拦截大部分重复消息这样既能减少对幂等表的写入压力又不会因为 Redis 故障导致完全失去保护。5. 方案三用手动提交 offset 控制投递语义5.1 关闭自动提交并改为手动提交默认情况下Kafka 消费者会周期性自动提交 offset。自动提交有两个问题提交周期内即使业务处理失败offset 也可能已经提交导致消息丢失处理时间不稳定时自动提交的时点和业务完成时点不一致重启后会出现重复。生产环境中如果业务要求更精确的控制通常会把自动提交关掉改为手动提交。在 spring-kafka 中配置大致是这样spring: kafka: consumer: bootstrap-servers: localhost:9092 group-id: order-consumer-group enable-auto-commit: false auto-offset-reset: earliest key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer listener: type: batch ack-mode: manual_immediate配合手动 acknowledgeComponent public class OrderMessageConsumer { KafkaListener(topics order-topic, groupId order-consumer-group) public void onMessage(ConsumerRecordString, String record, Acknowledgment acknowledgment) { try { handleOrder(record.value()); acknowledgment.acknowledge(); } catch (Exception e) { log.error(order consume error, message{}, record.value(), e); // 根据业务决定是否抛出异常交给容器重试或进入人工补偿 } } }关键区别在于只有 handleOrder 执行成功后才调用 acknowledge() 提交 offset。如果业务异常offset 不提交下一次 poll 会重新拿到这条消息。5.2 手动提交的常见误区误区一把 acknowledge() 写在 try 块外面。这样一旦业务异常异常没有阻止提交offset 照常前进消息就丢了。误区二以为手动提交就百分百不重复。只要处理成功和提交成功之间存在时间窗口仍然可能重复。手动提交只是把“是否认为消费完成”的判定权交到业务手里并不能消除崩溃窗口。误区三手动提交失败后不处理。提交失败意味着 offset 还停在旧位置业务却已经执行成功。此时需要根据场景选择继续重试提交、记录补偿日志或者让消息进入死信队列人工处理。5.3 相关参数速查表参数默认值含义调小的影响调大的影响生产建议enable.auto.committrue是否自动提交 offset手动提交时代码复杂度增加几乎不用关闭自动提交改用手动提交auto.commit.interval.ms5000自动提交间隔提交更频繁触发 stale offset 异常概率上升提交间隔大重启后重复范围更大根据业务容忍度调整max.poll.records500一次 poll 返回消息数单批处理量小避免超时单批处理量大容易超过处理时间上限按单条消息耗时和 max.poll.interval.ms 估算max.poll.interval.ms300000两次 poll 最大间隔处理超时更容易触发重平衡消费者失联被更晚发现大于最慢一条消息的处理时间session.timeout.ms45000会话超时时间心跳抖动更容易被判死消费者异常后组内感知变慢需要结合网络抖动情况设置heartbeat.interval.ms3000心跳发送间隔心跳频繁资源消耗上升心跳相隔较长误判风险上升应小于 session.timeout.msauto.offset.resetlatest没有提交 offset 时从哪开始丢失历史消息大量重复消费按业务语义选择 earliest 或 latest不同版本的 Kafka 客户端默认值可能会有变化落地前要以当前项目使用的客户端版本文档为准。参数调优的目标不是追求某一个值而是让“批量拉取数量、单条消息耗时、最大处理间隔”三者匹配。6. 生产环境怎么排查重复消费6.1 先看消费者组状态和 Lag排查重复消费第一步不是改代码而是观察消费者组的实时状态。用 Kafka 自带的命令行工具kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group order-consumer-group --describe也可以使用常见的 Kafka 可视化管理工具查看消费者组状态但在服务器上排查问题时命令行是第一步。输出里需要重点看几列CURRENT-OFFSET消费者组当前提交的 offsetLOG-END-OFFSET分区最新写入位置LAG两者差值代表积压消息数CONSUMER-ID当前分配到该分区的消费者实例。如果 LAG 很小但业务日志里出现大量处理记录基本可以断定同一批消息被重复消费。此时要马上确认业务处理时是否用了幂等键。如果 GROUP-STATE 长期停在 PreparingRebalance 或 CompletingRebalance说明消费者组在频繁重平衡。结合消费者端日志里的 Rebalance 记录可以判断是哪类触发条件。6.2 从消费端日志定位具体原因消费端日志中以下关键字都是重点排查对象日志关键字问题方向CommitFailedException提交 offset 时消费者已被移出组通常是处理时间超时导致max.poll.interval.ms两次 poll 间隔过长消费者被判失联WakeupExceptionpoll 被唤醒通常与重平衡或停止逻辑有关Offset commit failed提交失败需要看失败原因并检查补偿逻辑Rebalance组内成员变化或心跳超时session.timeout.ms心跳发送或网络异常看到上述关键字后把消费者跑一遍时间线poll 时间、业务开始处理时间、业务结束时间、ack 时间。如果业务结束时间与 ack 时间之间出现大量异常提交窗口就会变大。6.3 排查链路从现象倒推根因重复消费的排查建议按以下顺序进行业务侧是否已经有幂等记录幂等键是否唯一、覆盖面是否完整消费者组当前状态是 Stable 还是重平衡中消费者是否频繁 poll单条业务处理耗时是否超过 max.poll.interval.ms提交方式是否是手动提交ack 是否在业务成功之后提交失败时是否记录日志、是否走重试或人工补偿对比业务产生时间和重复消费时间判断重复来自重启、重平衡还是积压回放。按这个顺序排查大多数重复消费都能归到三类根因幂等缺失、提交时间窗口过大、消费者被重复拉入组导致分区重新分配。7. 不同环境下的落地建议与检查清单7.1 学习环境先跑通最小闭环学习阶段建议用 Docker 启动一个单节点测试环境docker run -d --name kafka-test -p 9092:9092 apache/kafka:latest示例中使用的是最新镜像实际环境建议指定稳定版本避免版本变化影响测试结果。单机测试时可以先不关注消费者数量和分区数的关系重点是验证消费者如何消费消息关闭自动提交后不调用 ack 会发生什么修改 group.id 后消费起点如何变化auto.offset.reset 设置为 earliest 和 latest 的区别。把上面几个点都验证一遍对重复消费的理解会比只看文章深很多。7.2 生产环境额外要补齐的事项生产环境不能只靠“手动提交 偶尔重启看效果”至少还要补齐这几样幂等设计数据库唯一键或 Redis 去重必须存在消费监控消费者组 Lag、重平衡次数、处理耗时都要纳入监控异常处理消费异常要区分“可重试”和“不可重试”不可重试的消息建议进入死信队列或人工补偿流程分区与消费者数量消费者实例数不要超过分区数否则多余实例闲置还会触发无意义重平衡配置外置化groupId、topic、消费超时时间等参数放到配置中心避免修改配置要重新发布发布与回滚消费者程序发布时要考虑优雅下线确保旧实例处理完再下线避免下线瞬间产生重平衡。7.3 发布前检查清单检查项检查结果消费者代码是否使用手动提交ack 是否在业务成功之后是 / 否消费逻辑是否有 try-catch异常是否可以正常回滚业务是 / 否幂等键是否能覆盖所有重复场景是 / 否幂等表和业务表是否在同一个事务内是 / 否group.id 是否与测试环境隔离是 / 否max.poll.interval.ms 是否大于最慢消费场景耗时是 / 否消费者数量是否与分区数匹配是 / 否是否已接入消费者组 Lag 监控是 / 否是否验证过重启后不会因 offset 未提交而大量重放是 / 否是否有消费失败后的补偿或死信机制是 / 否这份清单可以直接用于上线评审。每一项都做到重复消费仍然可能在极端故障下出现但它的影响会被限制在可控范围。Kafka 避免重复消费本质上不是“让 Kafka 不重复”而是“让重复消费变得无害”。面试时能把 at least once 语义、offset 提交窗口和三种方案讲清楚就已经超过大多数只背八股文的候选人。真正到生产环境建议先把幂等表或 Redis 去重落地再动手改手动提交参数和监控。只有先保证业务不重复执行后面的提交优化才有意义。