公司真题库

【阿里】Kafka 为什么出现重复数据,项目里怎么处理

91学AI·2026/7/20·7 阅读

考察点

这道题出自阿里数仓实习一面,前半问考原理后半问考实战。面试官想确认你知道重复不是 Kafka 的 bug 而是分布式系统 at-least-once 设计取向的必然结果,能从生产和消费两端把重复场景枚举出来;后半问"项目里怎么处理"是看实习生有没有真的接触过数据链路,能讲出具体的去重手段而不是背概念。追问常往"幂等生产者能完全解决吗"、"位移先提交还是先处理"、"hive 去重具体怎么写"走。

参考答案

生产者侧的重复

生产者发消息后等 broker 的 ack,如果网络抖动或者 broker 回应超时,生产者无法确定消息到底写没写进去。为了保证不丢,默认行为是重试——如果之前那条其实已经写成功了,重试就产生重复。这是分布式系统"两军问题"的典型体现:在不可靠网络上,"至少送达一次"和"最多送达一次"不可兼得,Kafka 默认选择不丢,代价就是可能重。

Kafka 0.11 之后有幂等生产者(enable.idempotence=true):broker 给生产者分配 PID,每个分区的消息带递增序号,重复到达的批次被拒绝。它能兜住"单个生产者会话内、单分区"的重试重复,但兜不住生产者重启(换了 PID)和跨分区场景。所以开了幂等也不等于天下太平,这一点要在面试里主动说,显得理解到位。

消费者侧的重复

消费侧的重复更常见,主要有两个来源。一个是位移提交时机:消费者处理完一批数据后、提交 offset 之前宕机,恢复后从上一次提交的 offset 重新拉取,这批已处理未提交的数据就重复消费了。反过来如果先提交 offset 再处理,宕机会导致丢失——所以"先提交"是 at-most-once,"后提交"是 at-least-once,语义的开关就在这个顺序上。

另一个是 rebalance:消费组里加机器、减机器、某个消费者心跳超时,都会触发分区重新分配。rebalance 瞬间如果某个分区从消费者 A 划给消费者 B,而 A 还有处理完未提交的数据,B 从旧 offset 接着拉,重复就产生了。rebalance 在生产上是高频事件(发布、扩缩容都会触发),这是消费侧重复的最大来源。

项目里的处理:分层设防

真实项目里不会指望单一手段,而是链路各环节各管一段。

入口侧开幂等生产者,成本为零(一个配置),先把生产重试的重复挡掉。

消费侧用 Flink 的场景(实时链路):位移不自动提交,而是纳入 checkpoint 状态,配合 sink 幂等写出做端到端。比如写 MySQL/Hologres 用业务主键 upsert(insert on duplicate key update),写 Elasticsearch 指定文档 id 覆盖——重复消费多少遍,外部系统里效果都等价于一次。这是"at-least-once 传输 + 幂等写入 = 效果上 exactly-once"的标准打法,成本低、够可靠,是大多数项目的真实选择。

离线侧(日志落 Hive 的链路):重复数据跟着 Kafka 消费落进 ODS,在 ODS 到 DWD 的清洗层去重。去重的关键是有一个业务主键——埋点日志一般有 event_id(端上生成的 UUID),业务库 binlog 有主键加操作时间。Hive 里的标准写法是 row_number 开窗:

INSERT OVERWRITE TABLE dwd_log_event_di PARTITION (dt='2026-07-01')
SELECT event_id, user_id, event_type, event_time, ...
FROM (
  SELECT *,
         row_number() OVER (PARTITION BY event_id ORDER BY event_time DESC) AS rn
  FROM ods_log_event_di
  WHERE dt='2026-07-01'
) t
WHERE rn = 1;

没有 event_id 的历史遗留埋点,退而求其次用"用户+事件+时间戳+关键属性"拼接近似主键,承认有极小误判率。

兜底:监控与对账

工程上最后一道防线是知道重复率。日常监控"ODS 行数 vs DWD 去重后行数"的比率,超过阈值(比如千分之五)告警——重复率突然飙升往往意味着上游出了问题:消费组频繁 rebalance、生产者超时参数不合理、或者有人重启了任务。另外重要报表链路做 T+1 对账,和源端(业务库)按主键比对数量和金额,把"去重失效"当成数据质量事故来管。面试里提到监控和对账,说明你理解去重不是写完 SQL 就结束的事。

一句话收拢

重复的根因是分布式系统只能保证 at-least-once,处理思路不是"消灭重复"而是"让重复无害":能幂等的环节幂等,不能幂等的环节用主键去重,最后靠监控兜底。

可能的追问

  • 幂等生产者有什么代价和局限?broker 要为每个 PID 维护序号状态;只保单会话单分区,重启、跨分区不保证;吞吐略有损耗。所以严格场景要开事务。
  • Flink 消费 Kafka 时位移到底什么时候提交?开启 checkpoint 时位移写入 checkpoint 状态,不向 Kafka 提交(或提交仅作监控用),恢复时从 checkpoint 的位移重放,保证状态和位移一致。
  • 去重为什么放 DWD 而不是 ODS?ODS 要保持贴源、可回溯,任何加工都放下游;且开窗去重需要全量扫一个分区,在 ODS 做会拖慢接入时效。

评论 (0)

暂无评论,快来抢沙发吧!

91学AI

© 2026 91学AI · 按岗位学 AI 与大数据. All rights reserved.