Kafka 面试题里,“如何避免重复消费”是出现频率极高的一道题。多数候选人能背出“消费者做幂等”“开启手动提交”这种结论,但被追问“重复消息到底从哪条路径产生”“offset 提交失败之后会发生什么”“幂等方案在并发场景下会不会失效”时,往往答不完整。这次我们把这题拆开讲透:先用本地 Kafka 复现一次真实重复消费,再动手实现 Redis 去重、数据库唯一键、状态机三种幂等方案,最后总结出一套可以直接用于面试回答和工程落地的处理思路。
这篇文章不是纯概念讲解。我们会给出完整的 Docker Compose 启动参数、Spring Boot 消费端配置、手动提交示例、批量消费去重逻辑,以及消费组 Lag 排查命令。读完你可以照着跑一遍本地实验,把 Kafka 重复消费的产生机制、解决思路、代码细节全部串起来。
1. 核心能力速览
先把这道面试题涉及的知识点整理成速览表,方便你快速判断自己哪块有缺口。
| 维度 | 说明 |
|---|---|
| 问题本质 | Kafka 消费端默认是至少一次(at least once)语义,重复消费是可能出现的现象 |
| 重复消息三个来源 | 生产者重试、消费端 Rebalance、Offset 提交时机不对 |
| 面试推荐回答层级 | 先讲重复是怎么产生的,再讲消费端幂等,最后讲哪些参数能降低重复概率 |
| 核心代码能力 | Spring Boot + spring-kafka 手动提交,业务处理成功后再提交 Offset |
| 幂等方案 | Redis SETNX 去重、数据库唯一键、消息状态表、状态机更新 |
| 能降低但不能根治的配置 | enable.auto.commit=false、max.poll.interval.ms、max.poll.records、auto.offset.reset |
| Kafka 完全干掉重复的边界 | 仅限 Kafka 到 Kafka 的单系统 exactly once;跨数据库、跨接口做不到 |
| 适合场景 | 订单、支付、积分、通知、数据同步、流处理任务 |
| 不适合场景 | 强依赖跨系统事务的业务,不能靠 Kafka 单独保证不重复 |
从表格能看到,这道题的核心矛盾在于:Kafka 本身并不承诺“一条消息只被处理一次”,它只承诺“分区内消息有序”。所以面试官真正想考察的,是你有没有在消费端做好数据一致性的兜底设计。
2. 重复消费问题全景
先明确一个事实:Kafka 的重复消费不是 bug,而是至少一次投递语义下的自然结果。生产者也一样,Kafka 提供了幂等生产者来避免分区内重复写入,但消费端没有等价机制,因为消费端要对接的外部系统 Kafka 管不了。
重复消费主要有三条路径。
第一条路径是生产者重试。生产者发送消息时如果网络抖动、Leader 选举导致发送超时,可能发起重试。重试一旦成功,Broker 上就会存在多条相同业务数据的消息。这种情况在 Kafka 0.11 之后可以用幂等生产者缓解,也就是enable.idempotence=true,它通过序列号机制让 Broker 识别重复批次,只写入一次。
第二条路径是消费者 Rebalance。消费者在处理一批消息时,如果处理耗时超过了max.poll.interval.ms,或者心跳超时导致 Consumer 被判定为宕机,就会触发 Rebalance。Rebalance 会把这个消费者的分区重新分配给其他消费者,而新消费者会从上次提交的 Offset 开始消费。如果旧消费者还没来得及提交 Offset 就被踢出,分区重新分配后就会从头再读一遍已经处理过的消息。
第三条路径是 Offset 提交时机。这是实际项目里最常见的出错点。很多新手配置enable.auto.commit=true,默认每隔几秒自动提交一次 Offset。如果消息已经拉取到本地,业务逻辑正在处理,自动提交还没发生,这时候进程崩了或者 Rebalance 了,重启后就会从旧 Offset 继续消费。表面上看是“消费者没消费完”,实际上是一条消息已经被处理,只是提交动作没跟上。
画成时间线就是:poll 一批消息 -> 业务处理 -> 还没 commit -> 崩溃/Rebalance -> 重新 poll -> 重复处理。
理解这三条路径之后,面试回答的第一层就有了:重复消费不能完全消灭,只能在消费端通过幂等和合理的提交策略来消除重复带来的影响。
3. 避免重复消费的主流方案
面试时不需要把所有方案全部铺开,重点讲清楚适合自己项目的那一种,然后补充说明其他方案的取舍。面试官更愿意听到“我为什么选它”而不是“我背了八个方案”。
主流方案可以分成三类。
第一类,消费端幂等。核心思想是让“重复执行”和“执行一次”产生相同结果。典型做法包括:
- 消息体中携带全局唯一业务 ID。
- 消费时用 Redis
SETNX判断是否处理过。 - 数据库表加唯一索引,插入冲突时直接忽略。
- 维护一张消息消费表,记录已处理的业务 ID。
第二类,调整消费端参数,降低重复消费概率。比如关闭自动提交、使用手动提交commitSync、把max.poll.interval.ms调大、把max.poll.records调小等。这种做法只能降低概率,不能根治,因为任何宕机和网络问题都可能导致提交失败。
第三类,使用 Kafka 事务。Kafka 本身提供事务 API,配合isolation.level=read_committed可以做到 Kafka 内部的 exactly once。但这只适用于“读 Kafka、处理、写回 Kafka”的场景,比如 Kafka Streams 或 Flink 写入 Kafka。一旦你的业务涉及数据库、第三方接口、下游服务,事务就覆盖不到,最终还是要靠业务幂等。
面试回答模板可以是:重复消费先分析来源,能通过参数优化的先优化,但必须在消费侧做幂等兜底,因为只有幂等才能在异常场景下保证数据一致。
4. 环境准备与前置条件
为了把方案落地验证一遍,我们需要在本地搭一套最小 Kafka 环境。推荐用 Docker Compose 启动 Kafka 3.5 以上版本,使用 KRaft 模式,不需要单独启动 ZooKeeper。
环境要求如下:
| 检查项 | 推荐配置 |
|---|---|
| 操作系统 | Linux、macOS、Windows 均可 |
| Docker | Docker 20.10 以上,Docker Compose V2 |
| JDK | JDK 8 或 JDK 11 以上 |
| 构建工具 | Maven 3.6 以上,或 Gradle |
| Spring Boot | 2.5 以上即可,这里以 Spring Boot 2.7 为例 |
| Kafka 版本 | 3.5 以上,KRaft 模式 |
| 端口 | 9092 用于 Kafka 服务,9000 用于 UI 界面(可选) |
| 磁盘空间 | 2GB 以上,Kafka 日志和 Docker 镜像会占用空间 |
如果你不想用 Docker,也可以直接下载 Kafka 二进制包,用kafka-server-start.sh启动。但 Docker Compose 更适合快速复现和清理环境,所以我们用 Docker。
另外建议准备一个 Kafka 可视化工具。桌面端可以用 Offset Explorer,也可以启动 Kafka UI 的 Docker 容器。可视化工具不是必须的,但排查消费组 Lag 和查看 Topic 分区时会方便很多。
5. 本地 Kafka 部署与 Spring Boot 消费者接入
先用 Docker Compose 启动 Kafka。这里使用 KRaft 模式,配置中不出现 ZooKeeper。
version: '3.8' services: kafka: image: bitnami/kafka:3.7 container_name: kafka-local ports: - "9092:9092" environment: - KAFKA_CFG_NODE_ID=1 - KAFKA_CFG_PROCESS_ROLES=broker,controller - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 - KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT - KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true volumes: - kafka-data:/bitnami/kafka volumes: kafka-data:在项目目录下保存为docker-compose.yml,然后执行:
docker compose up -d docker compose logs -f kafka看到类似Kafka Server started的日志就说明启动成功。
接着创建测试 Topic。先进入容器,或者直接在本机执行 Kafka 自带的脚本。如果没有二进制包,可以进容器执行:
docker exec -it kafka-local /opt/bitnami/kafka/bin/kafka-topics.sh \ --bootstrap-server localhost:9092 \ --create \ --topic order-events \ --partitions 3 \ --replication-factor 1然后创建一个 Spring Boot 项目,引入依赖:
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency>这里把 Redis 也引入进来,后面做幂等方案时直接用。
配置文件application.yml:
spring: application: name: kafka-repeat-consume-demo kafka: bootstrap-servers: localhost:9092 consumer: group-id: order-service auto-offset-reset: latest enable-auto-commit: false key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer properties: spring.json.trusted.packages: "*" listener: ack-mode: manual_immediate producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer注意这里的两个关键配置:
enable-auto-commit: false:关闭自动提交,改用手动提交。ack-mode: manual_immediate:消费者处理完消息后手动确认。
6. 模拟重复消费与问题定位
搭建好环境后,先写一个不带幂等的消费者,故意复现重复消费场景。这一步很有价值,面试时可以讲“我实际复现过”。
先写一个生产者,模拟发送订单消息:
@Service public class OrderProducer { private final KafkaTemplate<String, String> kafkaTemplate; public OrderProducer(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void send(String orderId) { // 消息体里放一个全局唯一业务ID String message = "{\"orderId\":\"" + orderId + "\",\"content\":\"create order\"}"; kafkaTemplate.send("order-events", orderId, message); } }再写一个普通消费者:
@Component public class OrderConsumer { private static final Logger log = LoggerFactory.getLogger(OrderConsumer.class); @KafkaListener(topics = "order-events", groupId = "order-service") public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) { log.info("收到消息: {}", record.value()); try { // 模拟业务处理耗时长,比如调用下游接口或者写数据库 Thread.sleep(1000); log.info("消息处理完成: {}", record.key()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { ack.acknowledge(); } } }这个消费者看起来已经使用了手动提交。但如果你把Thread.sleep时间改成明显超过max.poll.interval.ms的值,比如 6 分钟,就可以触发消费者被判定为失效,进而触发 Rebalance。Rebalance 发生后,这个消费者分到的分区被交给其他消费者实例,Offset 又未能提交,就会出现重复消费。
实际复现时不需要真的等 6 分钟。可以把 Topic 设置成多个分区,同时启动两个消费者实例,然后故意让其中一个消费者在处理消息时阻塞一段时间,再观察日志里是否出现同一个orderId被两次消费。下面命令用于查看消费组状态和 Lag:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group order-service这个命令的输出会展示每个分区的current-offset和log-end-offset。正常情况下两者差值趋近于 0,如果重复消费导致 Offset 回退,你会在日志中看到部分分区重复拉取数据。
到这里,我们已经复现了重复消费的第一现场。接下来解决它。
7. 幂等消费方案实战
前面说过,避免重复消费的根治法则是“消费端幂等”。下面给出三个可以直接上手的方案。
7.1 基于 Redis SETNX 的消息去重
每次消费前,用消息中的全局唯一业务 ID 作为 Key,Redis 的setIfAbsent作为判重依据。如果返回值是true,说明这条消息第一次处理,可以继续执行业务;如果返回false,说明已经处理过,直接确认消息并返回。
@Component public class OrderIdempotentConsumer { private static final Logger log = LoggerFactory.getLogger(OrderIdempotentConsumer.class); private final StringRedisTemplate redisTemplate; public OrderIdempotentConsumer(StringRedisTemplate redisTemplate) { this.redisTemplate = redisTemplate; } @KafkaListener(topics = "order-events", groupId = "order-service") public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) { String message = record.value(); String orderId = parseOrderId(message); // 设置去重Key,过期时间要大于业务最大处理时间 String key = "biz:order:" + orderId; Boolean firstConsume = redisTemplate.opsForValue() .setIfAbsent(key, "1", Duration.ofHours(1)); if (!Boolean.TRUE.equals(firstConsume)) { log.info("重复消息,直接确认并返回, orderId={}", orderId); ack.acknowledge(); return; } try { // 这里放真正的业务逻辑 handleOrder(message); log.info("订单处理成功, orderId={}", orderId); } catch (Exception e) { // 业务处理失败,需要删除Redis Key,否则后续重试会被幂等挡住 redisTemplate.delete(key); throw e; } finally { ack.acknowledge(); } } }这个方案的注意点有两个。第一是 Redis Key 的过期时间必须大于业务逻辑的最长处理时间,否则消息处理还没结束,Key 就过期了,重试消息会再次进来。第二是业务处理失败时不能保留 Key,要删掉,否则这条消息会一直被幂等拦截,永远不会被正常处理。
setIfAbsent是原子操作,多线程场景下不会出现两个消费者同时拿到true的情况。对于单分区单消费者场景,这个方案已经足够;即使因为 Rebalance 导致多实例并发,也能正确判重。
7.2 基于数据库唯一键的幂等方案
如果项目里没有 Redis,或者希望把幂等记录持久化,可以使用数据库唯一键。核心思路是建一张消息消费记录表,业务 ID 作为唯一索引,插入成功则继续处理,插入冲突说明重复。
先建表:
CREATE TABLE message_consume_log ( id BIGINT AUTO_INCREMENT PRIMARY KEY, biz_key VARCHAR(128) NOT NULL, consume_time DATETIME NOT NULL, UNIQUE KEY uk_biz_key (biz_key) ) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4;然后在消费者中插入记录,捕获唯一键冲突:
@Component public class OrderDbIdempotentConsumer { private static final Logger log = LoggerFactory.getLogger(OrderDbIdempotentConsumer.class); private final JdbcTemplate jdbcTemplate; public OrderDbIdempotentConsumer(JdbcTemplate jdbcTemplate) { this.jdbcTemplate = jdbcTemplate; } @KafkaListener(topics = "order-events", groupId = "order-service") public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) { String orderId = parseOrderId(record.value()); try { jdbcTemplate.update( "INSERT INTO message_consume_log(biz_key, consume_time) VALUES(?, ?)", orderId, LocalDateTime.now() ); } catch (DuplicateKeyException e) { log.info("重复消息,直接确认并返回, orderId={}", orderId); ack.acknowledge(); return; } try { handleOrder(record.value()); log.info("订单处理成功, orderId={}", orderId); } finally { ack.acknowledge(); } } }这里的关键是事务边界。如果handleOrder里有写库操作,最好把“插入消费记录”和“业务写库”放在同一个数据库事务中,否则可能出现消费记录写入成功、业务处理失败的情况。此时消息虽然被重复消费,但业务数据已经写库,最终消费者抛异常后 Kafka 重试,会直接走进DuplicateKeyException分支,看起来“消息被吃掉”了,实际业务数据反而没问题。
更严格的做法是把消费记录表的状态作为业务状态的一部分,用乐观锁控制更新。
7.3 基于状态机的幂等控制
另一种常见方案是维护业务单据状态。比如订单表中有status字段,处理消息时只执行合法的状态流转。订单状态从CREATED流转到PAID,如果消息重复投递,第二次消费时订单已经是PAID,就不会再做一次支付操作。
@Component public class OrderStateConsumer { @KafkaListener(topics = "order-events", groupId = "order-service") public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) { String orderId = parseOrderId(record.value()); // 只处理状态为CREATED的订单 int updated = orderMapper.casUpdateStatus(orderId, "CREATED", "PAID"); if (updated == 0) { log.info("订单状态已更新,忽略重复消息, orderId={}", orderId); ack.acknowledge(); return; } // 继续后续业务 handleOrder(orderId); ack.acknowledge(); } }casUpdateStatus在使用时通常是:
UPDATE t_order SET status = 'PAID' WHERE order_id = #{orderId} AND status = '#{expectStatus}'这个方案的好处是不需要额外的 Redis 或消费记录表,天然适合存在业务状态的场景。缺点是只能覆盖“有明确状态流转”的业务,对于没有状态的通用消息不太合适。
三种方案里,面试时优先讲第二种或第三种,因为它们设计得更有深度,能体现你对并发和事务边界有考虑。
8. Kafka 自身配置能做与不能做
在面试和技术方案评审中,经常会有人问:“能不能通过调整 Kafka 参数避免重复消费?”答案要分两层。
能做的是降低重复概率。关掉自动提交,改用commitSync,确保业务执行成功后再提交 Offset;调大max.poll.interval.ms,避免业务处理慢导致消费者被踢出;调小max.poll.records,减少单次拉取的消息量,降低处理超时风险。这些参数都是从“减少 Rebalance 和 Offset 滞后”的角度来降低重复消费频率,但不能消除网络异常、进程崩溃带来的重复。
不能做的是保证业务侧 exactly once。enable.idempotence=true只解决生产者到 Broker 的重复写入;isolation.level=read_committed只解决消费者读取事务消息的隔离级别;Kafka 事务只能保证 Kafka 内部生产者多个分区写入的原子性。一旦消息要从 Kafka 出去写入 MySQL、Redis、第三方 API,Kafka 就没有能力控制“外部动作执行一次”这件事。
所以,最稳妥的工程结论是:Kafka 配置优化 + 消费端幂等必须组合使用。前者降低重复概率,后者兜底保证最终一致。
9. 批量任务消费与失败重试
面试中还会顺带问一个延伸问题:“如果一条消息处理失败,消费组会重试吗?”这涉及到批量消费的设计。
Spring Kafka 支持批量监听,通过batch=true让消费者一次性拉取多条消息:
@Component public class BatchOrderConsumer { private final StringRedisTemplate redisTemplate; public BatchOrderConsumer(StringRedisTemplate redisTemplate) { this.redisTemplate = redisTemplate; } @KafkaListener(topics = "order-events", groupId = "order-service", batch = "true") public void onBatch(List<ConsumerRecord<String, String>> records, Acknowledgment ack) { // 先做批内去重,用内存Map覆盖同一批内的重复业务ID Map<String, ConsumerRecord<String, String>> distinctRecords = new HashMap<>(); for (ConsumerRecord<String, String> record : records) { distinctRecords.putIfAbsent(record.key(), record); } for (ConsumerRecord<String, String> record : distinctRecords.values()) { String orderId = record.key(); String key = "biz:order:" + orderId; Boolean firstConsume = redisTemplate.opsForValue() .setIfAbsent(key, "1", Duration.ofHours(1)); if (!Boolean.TRUE.equals(firstConsume)) { continue; } try { handleOrder(record.value()); } catch (Exception e) { log.error("订单处理失败, orderId={}", orderId, e); // 失败的消息记录到重试队列,或者发送到死信Topic redisTemplate.delete(key); } } ack.acknowledge(); } }批量消费需要注意几个点:
- 批内去重是必要的。Kafka 在同一批记录中可能就包含重复消息,先用内存 Map 过滤一遍,能减少无效的 Redis 或数据库访问。
- 批量模式下提交时机要正确。一个常见错误是等所有消息处理完才提交 Offset,如果中间有一条消息处理失败,整个批次都需要重试,会造成更多重复。更合理的做法是对单条失败的消息记录重试队列,剩余消息正常处理。
- 失败重试建议走独立 Topic 或本地定时任务补偿,而不是无限循环重试一条消息,否则会阻塞消费线程。
如果你在设计订单、支付这类消息消费链路,批量消费时一定要设计失败隔离机制,否则一条坏消息可能拖垮整批消息的处理进度。
10. 常见问题与排查方法
实际项目里,重复消费问题往往混杂着参数配置问题、网络问题和业务逻辑问题。下面整理了一份排查表,覆盖 Kafka 重复消费场景里最常见的几个问题。
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 消息被重复消费 | Offset 未及时提交或发生 Rebalance | 查看消费者日志,对比消费时间点;使用 kafka-consumer-groups.sh 查看 offset | 关闭自动提交,手动 ack;消费端加幂等 |
| Rebalance 频繁发生 | 处理耗时超过 max.poll.interval.ms,或心跳超时 | 查看日志中 rebalance 事件,观察 poll 间隔 | 调大 max.poll.interval.ms,调小 max.poll.records |
| 消费组没有拉到新消息 | group.id 不一致,或订阅了错误 Topic | 查看消费组描述和 Topic 分区列表 | 核对 group.id 和 topic 名称 |
| 启动时消费到大量历史消息 | auto.offset.reset 配置为 earliest | 查看消费者起始 offset 日志 | 按业务需求改为 latest,或手动重置 offset |
| 手动提交 Offset 抛异常 | commitSync 时网络波动或分组已失效 | 查看异常堆栈和 Consumer 状态 | 记录失败消息,配合幂等消费;使用 commitAsync + 回调 |
| Docker 启动 Kafka 后客户端报 Error while fetching metadata | 容器内监听地址和宿主机访问地址不一致 | 检查 KAFKA_ADVERTISED_LISTENERS 配置 | 设置 PLAINTEXT://localhost:9092 并映射宿主机端口 |
| 消费延迟很高,Lag 持续增长 | 单条消息处理太慢,或分区消费者数量不够 | kafka-consumer-groups.sh --describe 查看 Lag | 调小批量大小;提高消费者并发;优化下游性能 |
| Redis 幂等 Key 过期导致重复处理 | Key 过期时间设置过短 | 检查业务最慢处理耗时,查看 Redis TTL | 延长过期时间,或改成数据库唯一键方案 |
排查重复消费问题时,不要上来就改代码。先看几个关键信息:消费组当前 Lag、消费者实例数量、最近一次 Rebalance 时间、Offset 是否连续更新。多数情况下,问题都出在“处理慢导致 Rebalance”或者“手动提交时机不对”这两类原因上。
11. 最佳实践与使用建议
面试回答只需要讲清楚原理,但项目落地时,下面这些建议能帮你少踩坑。
第一,先判断业务是否能容忍少量重复。如果是纯日志上报、热度统计这类场景,重复消费影响不大,不需要特意做幂等;如果是订单、支付、转账、发放优惠券,必须做幂等,而且要选择可靠的幂等方案。
第二,手动提交 Offset 和业务成功必须严格绑定。不要在业务代码执行之前调用ack.acknowledge(),否则业务失败后消息不会被重试,造成数据丢失。正确顺序是:业务处理成功 -> 提交 Offset;业务处理失败 -> 不提交,让 Kafka 在当前消息上重试。
第三,幂等逻辑要考虑“处理中”的并发场景。如果消息真的被两个消费者实例同时处理,RedissetIfAbsent和数据库唯一键都能挡住;但如果先查 Redis 再插入,就会存在竞态窗口。尽量使用原子判等手段,不要先查询再判断。
第四,建立监控和告警。最基础的是消费组 Lag 监控,Lag 持续增长说明消费速度跟不上生产速度。还可以记录每条消息从生产到消费的耗时,超过阈值告警,快速发现消费者卡顿和 Rebalance 异常。
第五,合规和数据安全边界。消息数据中可能包含用户 ID、订单金额、联系方式等敏感信息,Kafka 日志、消费日志、幂等表里的数据要按公司数据安全规范处理,测试环境不要使用真实全量数据,不要将生产 Topic 数据导出到未授权环境。
12. 总结与下一步
Kafka 重复消费问题的本质,是 at least once 语义下的必然可能。面试回答要分三层:
第一层,重复消息的三条来源:生产者重试、Rebalance、Offset 提交时机不对。 第二层,常用规避手段:手动提交 Offset、合理设置max.poll.records、max.poll.interval.ms。 第三层,兜底方案:消费端幂等,具体实现可以是 Redis 去重、数据库唯一键或者状态机。
最容易踩的坑是把 Kafka 的 exactly once 理解成万能方案。实际项目中,跨系统的一致性只能靠业务幂等来保证。
下一步,建议你照着文章内容在本地起一套 Kafka,先跑一个普通消费者,故意把处理时间调长触发 Rebalance,观察重复日志;再加一个幂等方案,对比重复消息被拦截的效果。这个过程做完,这道面试题基本就稳了,同时你也能把消费端参数、手动提交机制、批量消费的问题都串起来理解。