
数据重复这个问题表面上看就是个“删掉多余的”这么简单的事但真正在项目里做起来你会发现它牵扯到的不仅仅是写一行drop_duplicates或者DELETE FROM ... WHERE ...而是从一开始就该想清楚的脏数据治理命题。我这些年经手过的项目里几乎每一个跟数据沾边的系统无论大小最后都得专门腾出精力来跟重复数据做斗争。今天这篇就把数据预处理阶段的重复问题彻底聊透从根因到检测从策略到实操再到各种防不胜防的坑一次性讲明白。1. 一个被低估的问题数据重复为什么值得专门处理很多人觉得数据重复不就是脏数据里最常见的一种嘛出现了删掉就行。但实际上重复数据的影响范围远比你想象的要大而且它产生的根源也特别复杂。如果不把来龙去脉搞清楚很容易出现今天删完明天又重复或者删完之后统计数据反而对不上的尴尬局面。1.1 重复数据从哪里来先说说重复数据的来源。我总结了一下绝大多数场景下无非是下面这几类业务系统重复提交用户手抖连续点了两次提交按钮或者前端做了重试机制但没做幂等控制订单表里就会多出两条一模一样的数据。这类情况在电商、支付、报名系统里最常见。数据同步与合并从多个数据源采集过来的数据天然就存在“同一实体对应多条记录”的情况。比如用户表里同一台手机号在不同的子系统里各存了一条合并到数仓的时候就产生了重复。日志与消息重复在大数据链路里上游的日志采集、消息队列投递为了保证不丢数据往往采用至少一次at least once的语义这就导致下游消费时天然会收到重复消息。Kafka 的重复消费问题就是这么来的。批处理作业的重跑定时调度任务因为超时或失败被重复执行而写入逻辑又没有做幂等处理跑一遍就多一份全量数据。手工操作与历史遗留运维或者业务人员手动修改数据、导入导出时不小心把文件内容复制了两遍还有老系统里没有主键或者主键设计不合理的历史包袱。这几种来源本身没有高低之分但它们决定了你需要用什么样的策略去处理。如果是业务系统重复提交你更多要考虑在入口处做拦截如果是消息队列的重复消费那就要靠下游的幂等设计来兜底如果是多源合并那就要在 ETL 过程里按业务规则去重。数据预处理阶段的去重本质上就是在为上层的数据分析、模型训练和业务决策做第一道防线。1.2 数据重复造成的影响统计结果失真一个用户如果被统计了两次那么“活跃用户数”“订单总量”“用户平均消费”这些核心指标全部都是错的。如果这个错误指标恰好被用到了报表或者 KPI 考核上那问题就大了。模型训练过拟合或偏差在机器学习任务里重复样本会放大某些模式在训练集中的权重导致模型对重复出现的样本过拟合而对真实分布欠拟合。尤其在分类问题中重复数据还可能造成严重的类别不平衡。占用存储与计算资源在大数据环境下存储和计算成本是按量计算的。海量重复数据意味着你在为无效信息买单跑任务的时间也被白白拉长。下游任务直接报错比如在数据库里唯一索引冲突导致写入失败或者两张表关联时因为重复记录导致结果膨胀一对多变成多对多这类问题排查起来非常费劲。1.3 先定义清楚什么样的记录算重复这是整个去重流程里最容易出问题的一步。不同场景下“重复”的判定标准完全不一样完全重复所有字段的值完全一致。这种最简单多见于日志数据重复采集或文件重复导入。业务键重复某个或某几个关键业务字段一样但其他字段可能不同。比如同一订单号出现了两条记录一条是待支付状态、一条是已支付状态这算重复吗要看你分析的粒度。如果分析的是订单数这就是重复如果分析的是订单状态流转这可能不是重复。相似重复数据本身存在格式差异比如“张三”和“张 三”或者“北京市朝阳区”和“北京朝阳区”需要做相似度匹配才能识别。这类是数据清洗的高级话题本文先不展开重点说清楚前两类。在写任何去重代码之前一定要跟业务方确认去重的键和去重的语义。否则后面做的一切都可能是在给数据“做手术”但没有对症下药。2. 重复数据检测先搞清楚哪些算重复再动手理解了“什么是重复”之后才能谈“怎么检测”。这一步的目的不是为了直接去重而是为了摸清数据里到底有多少重复、集中在哪些字段、是什么样的重复形态。搞清楚这些你才知道接下来的去重方案要不要做成定时任务、要不要引入人工审核机制。2.1 快速摸排用 SQL 或 pandas 看一眼重复规模不管数据是在数据库里还是在文件里第一件事永远是先查“重复到底有多少”。在 SQL 里可以这样SELECT 业务键, COUNT(*) AS cnt FROM 表名 GROUP BY 业务键 HAVING COUNT(*) 1 ORDER BY cnt DESC LIMIT 20;这组语句能快速找出重复次数最多的前 20 个业务键。先看数据再写逻辑是这个环节的黄金法则千万别跳过这一步直接就去重不然你连自己删掉了什么都不知道。在 Python 里用 pandas 做初步探查也很直接import pandas as pd df pd.read_csv(raw_data.csv) # 全字段重复统计 full_duplicated df.duplicated().sum() print(f全字段完全重复的行数: {full_duplicated}) # 指定业务键重复统计 key_duplicated df.duplicated(subset[user_id, order_id]).sum() print(f按 user_id order_id 判断的重复行数: {key_duplicated}) # 看重复样本长什么样 duplicate_only df[df.duplicated(subset[user_id, order_id], keepFalse)] print(duplicate_only.sort_values(by[user_id, order_id]).head(20))这里有个细节值得说一下keepFalse表示把所有重复行的全部记录都标出来而不只是保留其中一条或把重复的扔掉一部分。这一步对人工核验特别有用。2.2 全字段完全重复的检测全字段完全重复的检测在技术上是最简单的因为它的判断标准就是所有列的值一模一样。在 pandas 里直接用duplicated()默认参数就行。但要注意全字段完全重复的情况在实际业务数据里其实占比不高因为大多数业务表都会有时间戳、自增 ID、状态字段等这些字段稍微有点不一样整行就不是完全重复了。所以在实践中全字段完全重复往往出现在日志采集、埋点数据以及通过 API 批量导入导出的场景里。用df.duplicated().sum()拿到重复行数之后接着就要做定位和抽样看看重复数据是连着的还是散落的。# 找出所有完全重复行的索引 dup_mask df.duplicated(keepFalse) dup_rows df[dup_mask].sort_values(bydf.columns.tolist()).head(50) print(dup_rows)排序之后一眼就能看出重复数据的分布规律如果多条一模一样的记录是紧挨着的那极有可能是一次导入操作把同一份数据写入了多次如果重复行之间夹杂着其他记录那可能是业务上反复提交。2.3 基于业务键的部分重复检测业务键重复比完全重复要常见得多但它的判定也更依赖你对业务的理解。常见做法是先找出一组能够唯一标识一条业务记录的字段组合然后基于这些字段去做检测。比如用户交易场景里user_id order_id往往能定位到一条唯一的记录用户画像场景里user_id或手机号可能是唯一的内容分析场景里文档ID或标题 发布时间可能是唯一的。business_key [customer_id, transaction_date] df[is_duplicate] df.duplicated(subsetbusiness_key, keepFalse) # 看看哪些业务键重复次数最多 dup_count df.groupby(business_key).size().reset_index(namecount) dup_count dup_count[dup_count[count] 1].sort_values(count, ascendingFalse) print(dup_count.head(20))但这还没完。业务键重复经常伴随一个问题你拿来做去重依据的字段本身可能就存在脏数据比如空值、格式不一致。所以在检测的同时最好先做一下字段质量探查看看这些字段的非空率、唯一值数量、格式分布。如果某个业务键有大量空值那空值记录就会被单独归为一组可能造成误杀或漏杀。2.4 判断保留哪一条不是随便留一条就行的检测完重复后下一个问题就是如果同一组重复记录里有几条数据内容不一致应该保留哪一条很多人图省事直接keepfirst但这样做往往是有隐患的。举例来说一个订单表里有这么两条记录订单号状态金额更新时间A001待支付1002024-01-01 10:00:00A001已支付1002024-01-01 10:30:00如果你按订单号去重并简单保留第一条那保留的就是“待支付”状态的订单。可实际上后一条记录才是这个订单的最新状态。所以“保留哪条”在很多时候不是随便选的它取决于你想让这份数据表达什么以及哪个时间戳/状态最能代表当前真实事实。比较稳妥的做法是在去重前先按时间倒序排序然后再执行去重把最新的一条留下来。或者更精细一点针对同组重复记录把某些字段合并成聚合值比如取最大时间、最早时间、状态拼接等。# 按更新时间倒序排序后每个订单保留最新一条记录 df_sorted df.sort_values( by[订单号, 更新时间], ascending[True, False] ) df_dedup df_sorted.drop_duplicates(subset[订单号], keepfirst)这段逻辑看起来只有几行但它表达了一个非常重要的思想去重不是简单地丢弃数据而是在信息不完整的情况下选一条对后续分析最有意义的记录作为代表。3. 不同场景下的去重方案与选型去重的实现方案没有银弹关键看数据存在哪里、数据量大不大、是批量处理还是实时处理。这一章把我个人在各种场景下用过的方案和选型思路整理出来。3.1 关系型数据库SQL Server / MySQL里去重关系型数据库里去重最经典的手段就是窗口函数ROW_NUMBER()。这个方案的核心逻辑是先按业务键分组再按某个排序规则给组内记录编号最后只保留编号为 1 的记录。WITH ranked AS ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY 订单号 ORDER BY 更新时间 DESC ) AS rn FROM 订单表 ) DELETE FROM 订单表 WHERE 订单号 IN ( SELECT 订单号 FROM ranked WHERE rn 1 );执行这个操作之前强烈建议先开一个事务并且最好是先查后删用 SELECT 确认被删的数据范围是对的再执行 DELETE。还有一种场景是表里没有 ID、所有字段都一样的完全重复记录。这种在 SQL Server 里可以直接用DISTINCT查出去重后的结果然后重建表SELECT DISTINCT * INTO 订单表_去重 FROM 订单表; -- 确认数据没问题后重命名表 EXEC sp_rename 订单表, 订单表_带重复; EXEC sp_rename 订单表_去重, 订单表;但在生产环境里把表直接改名重建也算是个危险操作最好是在维护窗口做同时要停掉对这张表的写入服务。3.2 用 Python / pandas 处理批量文件数据当数据量大到数据库本身扛不住或者数据是以文件形式存在时pandas 依然是数据预处理的首选。上面已经写了很多 pandas 片段这里补充一个完整的、可直接复制的处理流程import pandas as pd # 1. 读取数据 df pd.read_csv(raw_orders.csv, dtype{订单号: str, 用户ID: str}) # 2. 字段标准化去空格、统一大小写 df[订单号] df[订单号].str.strip() df[用户ID] df[用户ID].str.strip() # 3. 按业务键去重保留最新一条 df[更新时间] pd.to_datetime(df[更新时间]) df_sorted df.sort_values(更新时间, ascendingFalse) df_dedup df_sorted.drop_duplicates(subset[订单号], keepfirst) # 4. 去重效果统计 before_count len(df) after_count len(df_dedup) print(f去重前记录数: {before_count}) print(f去重后记录数: {after_count}) print(f重复记录数: {before_count - after_count}) # 5. 输出到新文件 df_dedup.to_csv(orders_dedup.csv, indexFalse, encodingutf-8-sig)需要注意的一个点是read_csv时指定dtype参数可以把订单号、手机号这类字段按字符串读入避免长数字被读成科学计数法这一点很多人踩过坑。比如一个 16 位的订单号如果不指定 dtype读完就变成1.23456e15后续做匹配全部失败。当文件超级大比如几十 GB时pandas 的read_csv会内存爆炸。这时候可以换成分块读取chunksize 分批去重 落盘合并的方案或者直接上 Spark / Flink。关于大数据量下怎么做后面专门用一节说。3.3 流式场景Kafka、Flink下的去重问题Kafka 这类消息队列在数据链路里极其常用但它本身默认不保证消息不重复。生产端重试、消费端在提交 Offset 之前宕机都会导致下游收到重复消息。处理思路一般有两个方向方向一下游幂等。让最终写入的目标系统自己承担去重职责比如在数据库表里建唯一索引重复写入时采用INSERT ... ON DUPLICATE KEY UPDATE或者MERGE语句打平重复记录。方向二前置去重。在流处理引擎里维护一个窗口状态对滑动窗口内的数据做去重。Flink 里有DISTINCT算子和KeyedProcessFunction可以自己实现去重逻辑但要注意状态大小和失效时间。// Flink 中通过状态去重伪代码示意 DataStreamOrder stream ...; stream .keyBy(order - order.getOrderId()) .process(new KeyedProcessFunctionString, Order, Order() { private ValueStateBoolean seenState; Override public void open(Configuration parameters) { seenState getRuntimeContext().getState( new ValueStateDescriptor(seen, Types.BOOLEAN) ); } Override public void processElement(Order order, Context ctx, CollectorOrder out) throws Exception { if (seenState.value() null) { seenState.update(true); out.collect(order); } // 如果已经见过该订单号直接丢弃 } });这段代码虽然简单但背后有一个隐藏问题状态seenState会随着 Key 不断增长而膨胀。所以需要引入状态 TTL比如设置 24 小时后自动过期定期清理冷数据。对于实时场景来说时间窗口内的去重比永久去重更实际。3.4 数仓与 ETL 里的去重设计到了数据仓库层面去重就不再只是删几条数据的事了而是要设计出一套机制从源头上防止重复数据流入下游。常见做法包括在 ODS 层做“增量 幂等写入”每次抽取数据前先 Checkpoint记录本次导入的数据范围避免重复导入。在 DWD 层建业务主键模型明细事实表在建模时明确业务主键并针对该主键做唯一性校验。使用拉链表 / 快照表来处理缓慢变化维用哈希值或业务序列号避免重复更新。调度作业里加上“数据质量监控”每次 ETL 跑完后自动执行重复计数 SQL如果重复率超过阈值就报警比出了问题再修好得多。说到底数仓层级的去重要比文件级和数据库级去重多考虑一个因素可回溯性和血缘关系。你今天去重砍掉的那些数据之后碰到 “ 数据对不上 ” 的排查别人需要能顺着血缘找到你做了去重这一步。4. 完整实操从数据摸底到去重落库的标准化流程前面讲了不少策略和片段这一章我按照实际项目的流程串起来写一套可以从头跑到尾的完整处理流程。以一份 CSV 格式的订单流水为例处理目标是把同一个订单号对应的多条记录压缩成一条保留最新状态并输出去重报告。4.1 步骤一环境准备与数据探查先确认 Python 环境里有没有 pandas没有就安装pip install pandas然后写一个探查脚本摸清数据结构import pandas as pd df pd.read_csv(orders.csv, dtype{订单号: str, 用户ID: str}) print(列名:, df.columns.tolist()) print(总行数:, len(df)) print(各列非空统计:) print(df.notna().sum()) print(数据类型:) print(df.dtypes)这一步能帮你确认几个关键问题文件编码是否正确、字段名是否需要重命名、时间字段是不是日期类型、业务键字段是否有空值。别小看这几步留着脏字段直接跑去重最后一定会出现误删或漏删。4.2 步骤二制定去重规则根据业务需求我在这里定三条规则一个订单号只保留一条记录。如果同一订单号出现多条记录保留更新时间最新的一条。如果更新时间也相同则保留最后一个字段拼接起来做二次排序后最大的一条兜底逻辑保证结果稳定。规则的优先级和排序字段一定要写清楚这样代码执行结果才是确定性的。否则你跑两次任务输出的结果可能不同。df[更新时间] pd.to_datetime(df[更新时间]) df[sort_key] df[更新时间].astype(int64) # 纳秒时间戳 df df.sort_values([订单号, sort_key], ascending[True, False]) df_dedup df.drop_duplicates(subset[订单号], keepfirst)这里我把时间列转换成纳秒级整数当排序键理由是日期类型在 pandas 里的比较可能因为时区问题产生歧义转成整数以后排序逻辑完全透明可控。4.3 步骤三执行去重并校验结果校验是整个流程里最容易被忽略但最重要的环节。去重结果对不对不能光看行数减少了多少还要做抽样确认。# 抽样校验检查订单号唯一性 assert df_dedup[订单号].is_unique, 去重后订单号仍存在重复 # 抽查一个已知重复订单确认保留的是最新那条 sample_order A10001 subset df[df[订单号] sample_order].sort_values(更新时间, ascendingFalse) print(原始记录) print(subset) print(去重保留) print(df_dedup[df_dedup[订单号] sample_order]) # 重复率统计 total len(df) dedup_count len(df_dedup) print(f原始记录数: {total}) print(f去重后记录数: {dedup_count}) print(f重复比例: {(total - dedup_count) / total * 100:.2f}%)校验完毕后导出数据。如果下游是数据库可以直接用pd.to_sql写入但要注意提前清空目标表避免追加模式下再次制造重复数据。# 导出到 CSV df_dedup df_dedup.drop(columns[sort_key]) df_dedup.to_csv(orders_dedup.csv, indexFalse, encodingutf-8-sig) # 或者写入 SQL Server from sqlalchemy import create_engine engine create_engine(mssqlpymssql://user:passwordhost/dbname) df_dedup.to_sql(orders_dedup, conengine, if_existsreplace, indexFalse)if_existsreplace这个参数特别值得强调它会在写入前把旧表删掉重建能确保每次跑出来的结果都是全量最新的不会在表里留下上次任务的残留记录。4.4 步骤四生成去重质量控制报告一个有经验的数据工程师做完去重一定会顺手生成一份简单的质量报告发给业务方或者留在调度日志里备查。报告内容不用花哨核心指标够了report { 数据源文件: orders.csv, 原始记录数: total, 去重后记录数: dedup_count, 删除重复记录数: total - dedup_count, 重复率: f{(total - dedup_count) / total * 100:.2f}%, 去重键: 订单号, 保留策略: 保留更新时间最新记录, 处理时间: pd.Timestamp.now().strftime(%Y-%m-%d %H:%M:%S) } report_df pd.DataFrame([report]) report_df.to_csv(dedup_report.csv, indexFalse, encodingutf-8-sig) print(report_df)这份报告看着简单但在后续排查数据问题时特别有用。比如业务方问“为什么订单表里某张单子不见了”你把报告里记录的“删除重复记录数”和具体删除逻辑给他们沟通效率会高很多。5. 大数量与分布式场景下的去重策略文件几万行、数据库几百万行的时候pandas 和 SQL 都能轻松搞定。但一旦数据量到了亿级别或者数据分布在多个节点上情况就完全不一样了。这一章聊聊大数据量下的去重策略。5.1 数据量超过内存时怎么去重当数据文件达到几十 GB 级别单机内存根本装不下的时候有几种路子可以走。第一种是分块处理加外部排序import pandas as pd chunk_iter pd.read_csv(huge_data.csv, chunksize1000000, dtype{id: str}) seen set() # 只保存见过的业务键集合 for chunk in chunk_iter: chunk chunk.sort_values(更新时间, ascendingFalse) dedup_chunk chunk[~chunk[订单号].isin(seen)] seen.update(dedup_chunk[订单号].tolist()) # 按需落盘 dedup_chunk.to_csv(dedup_part.csv, modea, headerFalse, indexFalse)这种方案的问题也很明显seen集合会随着业务键数量增长而膨胀如果业务键有几千万个内存一样扛不住。所以这只适用于“业务键基数不是极其夸张”的中量级场景。第二种是上 Spark。用 Spark 做去重基本是分布式计算里的标准操作了而且对单机内存没有要求因为数据和状态都是分布在不同 executor 上的// Scala / Spark 伪代码 import org.apache.spark.sql.functions._ val df spark.read.option(header, true).csv(hdfs://path/to/orders.csv) val dfDedup df .withColumn(ts, to_timestamp(col(更新时间))) .sort(col(ts).desc) .dropDuplicates(订单号) dfDedup.write.mode(overwrite).csv(hdfs://path/to/orders_dedup.csv)用 Spark 去重有个需要注意的点sort是全局排序还是分区排序会直接影响到dropDuplicates保留下来的记录。如果你需要保留“全局最新的一条”最好的做法是先做一次全局排序或者用窗口函数 过滤来实现并且在集群资源上预留足够的内存避免执行阶段出现 OOM。import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ val windowSpec Window.partitionBy(订单号).orderBy(col(ts).desc) val dfDedup df .withColumn(rn, row_number().over(windowSpec)) .filter(col(rn) 1) .drop(rn)row_number()这个窗口函数的逻辑和 SQL Server 里一模一样只是从单机搬到了分布式环境。5.2 分布式去重的常见坑分布式下去重的难点不在于算法而在于数据倾斜和状态管理。数据倾斜在去重场景里特别常见。比如某些“爆款订单号”或者“大客户ID”下挂了海量记录某个分区的数据量是其他分区的几百倍作业执行时间被这个分区严重拖慢。这个问题的缓解方法是把倾斜的键加盐后打散再做两阶段去重。这个方案写起来有一点复杂但理解它的核心思想就够了先加随机前缀把数据扩散到更多分区做完第一轮局部去重后再去掉前缀做全局去重。状态管理主要发生在流式去重场景。如果用 Flink 或 Kafka Streams 做去重要格外注意状态后端的选型和 TTL 设置。默认情况下状态是保存在 RocksDB 或内存中的不设置 TTL 的话状态会无限增长最终拖垮整个作业。一般建议对去重状态设置一个合理的保留时间比如“24小时内重复数据视为重复超过24小时视为新数据”。5.3 要不要用布隆过滤器做大规模去重布隆过滤器Bloom Filter在大规模去重里经常被提到它的优势是用了远小于原始数据的内存就能判断“一个值是否出现过”。但它有一个致命弱点存在误判率它只会把不存在的值误判成存在也就是有假阳性没有假阴性。在实际项目中我一般把布隆过滤器用在“快速淘汰”阶段。比如在消费 Kafka 数据时先用布隆过滤器判断订单号是否已经见过如果过滤器说“没见过”那可以直接放行如果过滤器说“见过”再去查真正去重表做二次确认。这样能过滤掉大部分明显重复的数据同时不会因为误判造成数据损失。布隆过滤器不能单独作为去重方案的最终判定必须配合精确判定一起用。这是一个很多新手容易犯的错我也踩过。6. 常见问题与排查技巧实录最后整理一份我在实际项目中处理数据重复问题时的常见问题速查表。这些问题不一定每一次都会遇到但遇到一个就够折腾半天的。6.1 去重键判断失误导致有效数据被删这是最严重的事故类型。有一次我做用户画像时用手机号作为唯一去重键结果发现有些用户手机号是空的多条不同用户的记录因为空值被聚到了同一组里去重后直接丢失了真正业务意义上的不同用户。这个问题的根源在于空值处理不到位。在确定去重键之前一定要排查业务键的空值比例。如果空值比例过高可以考虑先拆分数据有业务键的数据按业务键去重无业务键的数据用其他规则如所有字段拼接的哈希值去重最后再合并。6.2 去重后数据量不减反增听着很荒谬但我确实遇到过。原因出在 join 操作上你在做去重时如果先把两张表 join 在一起而 join 的键本身又不是唯一的那么去重操作有可能把一个本来唯一的记录膨胀成多条。所以去重应该尽量在单张表、单份数据内完成如果必须在 join 之后去重要先确认 join 不会产生一对多的膨胀。6.3 排序与去重的顺序反了在 pandas 的drop_duplicates里很多人以为先执行再排序也没关系。其实不行drop_duplicates(keepfirst)保留的是“在该列第一次出现”的那条记录这个“第一次”取决于 DataFrame 当前的顺序。如果不先排序你保留的就可能不是你想要的“最新一条”而是文件里的“原始第一条”。# 错误示例先 drop_duplicates 再排序 df_dedup_wrong df.drop_duplicates(subset[订单号], keepfirst) df_dedup_wrong df_dedup_wrong.sort_values(更新时间, ascendingFalse) # 正确示例先排序再 drop_duplicates df_dedup_right df.sort_values(更新时间, ascendingFalse).drop_duplicates(subset[订单号], keepfirst)两段代码看似只有顺序差别但结果完全不同。排序先行的逻辑是让目标记录排在最前面这样去重时它就会被保留下来。6.4 流式系统中重复消费的定位思路Kafka 消费者出现重复消费时第一反应不应该是改代码加去重而是先定位重复的阶段。用kafka-consumer-groups.sh查看消费者组的 Lag同时比对日志中的 Offset 提交记录确认是“消费了但没提交 Offset”还是“重复投递但消费端重复处理”。只有定位清楚之后才能选对处理方式。如果是 Offset 提交失败那修复提交逻辑就行如果是业务系统自身重试导致的重复请求那就得在业务接口层做幂等。6.5 加唯一索引之后数据写入报错在清理完重复数据后很多人会顺手加一个唯一索引来防止未来再产生重复。这个思路没问题但前提是你已经做完了足够的数据质量检查。因为一旦唯一索引加在本来就有重复数据的历史表上系统直接报错连写入都做不了。加索引的流程应该是先清重复 → 再校验唯一性 → 再加索引 → 最后做写入验证。顺序不能乱。场景常见现象排查思路全字段重复duplicated().sum() 数值异常高检查上游采集是否重复发送、导入时是否重复执行业务键重复按 key 分组 count1 的记录很多确认业务键选择是否合理考虑空值和格式不一致流式重复消费消费者重复处理同一条消息查看 Offset 提交、故障重启策略加幂等去重后误删数据量明显减少且关键记录消失回看去重规则检查排序字段和空值处理逻辑大文件去重内存爆pandas 直接 OOM改用 Spark、分块处理或数据库/数仓方案再补充两个小技巧。第一个在排查重复问题时不要只看行数要多看“唯一键数量”。第二个不管用什么方式去重跑完之后一定要做行数比对和抽样检查确保满足预期再去下游消费。这些细节看似繁琐但在真实项目里能帮你省掉后面无穷无尽的返工时间。我在实际项目里最深的感受是去重这件事真正难的从来不是那个 drop 的动作而是想明白要去重到什么粒度、保留什么版本、用什么机制防止它再次发生。把这三层想透了无论是几万行的小文件还是几亿行的分布式大表思路都是一样的。