搞定删除重复数据保留一条:从报错到源码解析的实战指南
搞定删除重复数据保留一条:从报错到源码解析的实战指南
配置环境就卡半天,是不是熟悉的感觉?跑个脚本删个重数据,结果环境依赖打架,SQL 语法报错,或者数据删了但主键冲突,那种抓狂感真让人想砸键盘。别急,今天咱们不整虚的,直接上手一个完整的实战项目,通过源码解析彻底搞懂删除重复数据保留一条的底层逻辑。
项目目标与痛点场景
在真实的生产环境中,数据重复是常态。无论是用户注册时的并发提交,还是数据迁移时的脚本失误,重复数据(Duplicate Data)都会导致统计偏差、业务逻辑错误甚至系统崩溃。
我们的目标很明确:构建一个可复现的 Python 工具,能够安全、高效地删除重复数据保留一条记录。这里的“保留一条”,指的是保留最新(MAX(id))或最旧(MIN(id))的那一条,其余标记删除或物理删除。
为什么需要源码解析?因为网上大多数教程只给 DELETE 语句,一旦遇到大表、高并发或者特定数据库(如 MySQL、PostgreSQL、Oracle)的差异,代码直接崩盘。我们需要深入到底层执行逻辑,理解事务隔离级别、索引选择以及锁机制,才能写出真正健壮的代码。
目录结构与依赖环境
为了避免“配置环境就卡半天”的噩梦,我们先定义一个最小化、可复现的项目结构。不要使用复杂的框架,就用最原生的 Python 标准库和数据库驱动。
project_structure/
├── requirements.txt
├── config.py # 数据库配置
├── db_utils.py # 数据库连接池与执行器
├── duplicate_cleaner.py # 核心清洗逻辑
├── test_data.py # 生成测试数据
└── main.py # 入口脚本requirements.txt 内容如下,保持精简:
pymysql==1.0.2
psycopg2-binary==2.9.5
loguru==0.7.0关键避坑点:很多新人喜欢用 SQLAlchemy 做 ORM,但在处理大规模数据清洗时,ORM 的抽象层会带来巨大的内存开销和性能损耗。对于删除重复数据保留一条这种底层操作,直接使用 SQL 或 DB-API 接口是更高效、更可控的选择。
config.py 中我们定义多数据库支持,通过配置切换:
import os# 环境变量覆盖默认配置,方便 CI/CD 部署
DB_CONFIG = {engine: os.getenv(DB_ENGINE, mysql), # mysql, postgreshost: os.getenv(DB_HOST, 127.0.0.1),port: int(os.getenv(DB_PORT, 3306)),user: os.getenv(DB_USER, root),password: os.getenv(DB_PASSWORD, password),database: os.getenv(DB_NAME, test_db),batch_size: int(os.getenv(BATCH_SIZE, 1000)) # 批量处理大小
}核心代码实现与源码解析
这是本文的核心部分。我们将分两步走:先写一个“能跑”的版本,再通过源码解析找出它的隐患,最后给出优化后的生产级代码。
1. 基础实现:朴素 SQL 的陷阱
很多博客教你这样写:
DELETE t1 FROM users t1
INNER JOIN users t2
WHERE t1.id t2.id
AND t1.email = t2.email;这招在 MySQL 小表上确实好用。但在 Postgres 里,语法完全不同;在 Oracle 里,DELETE 不能直接子查询同一张表。更致命的是,这种写法没有事务控制,没有进度反馈,一旦数据量上百万,连接超时直接中断,数据处于“半删除”状态,灾难就此发生。
2. 生产级实现:Python 封装与批量处理
我们使用 loguru 记录日志,使用批量处理(Batching)避免长事务。
# db_utils.py
import pymysql
import psycopg2
from contextlib import contextmanager
from loguru import logger
import config@contextmanager
def get_db_cursor():上下文管理器,确保连接和游标正确关闭conn = Nonecursor = Nonetry:if config.DB_CONFIG[engine] == mysql:conn = pymysql.connect(host=config.DB_CONFIG[host],port=config.DB_CONFIG[port],user=config.DB_CONFIG[user],password=config.DB_CONFIG[password],database=config.DB_CONFIG[database],charset='utf8mb4',autocommit=False # 关键:手动提交,控制事务)elif config.DB_CONFIG[engine] == postgres:conn = psycopg2.connect(host=config.DB_CONFIG[host],port=config.DB_CONFIG[port],user=config.DB_CONFIG[user],password=config.DB_CONFIG[password],dbname=config.DB_CONFIG[database],)else:raise ValueError(fUnsupported engine: {config.DB_CONFIG['engine']})cursor = conn.cursor()yield cursorexcept Exception as e:logger.error(fDatabase error: {e})if conn:conn.rollback()raisefinally:if cursor:cursor.close()if conn:conn.close()接下来是核心清洗逻辑 duplicate_cleaner.py。这里我们采用“查找重复组 - 保留最大 ID - 删除其余”的策略。
# duplicate_cleaner.py
import config
from db_utils import get_db_cursor
from loguru import loggerdef find_duplicate_groups(table_name: str, unique_columns: list):查找所有存在重复的分组返回: [(group_values, max_id, count), ...]columns_str = , .join(unique_columns)# 注意:这里使用 GROUP BY 和 HAVING,避免全表扫描sql = fSELECT {columns_str}, MAX(id) as max_id, COUNT(*) as cntFROM {table_name}GROUP BY {columns_str}HAVING COUNT(*) 1results = []with get_db_cursor() as cursor:cursor.execute(sql)for row in cursor.fetchall():# row 结构: (col1, col2, ..., max_id, count)group_vals = row[:-2]max_id = row[-2]cnt = row[-1]results.append((group_vals, max_id, cnt))return resultsdef delete_duplicates(table_name: str, unique_columns: list, batch_size: int = 1000):执行删除操作,保留每组中 ID 最大的一条# 1. 获取重复组dup_groups = find_duplicate_groups(table_name, unique_columns)if not dup_groups:logger.info(No duplicates found.)return 0logger.info(fFound {len(dup_groups)} duplicate groups. Starting cleanup...)total_deleted = 0columns_str = , .join(unique_columns)# 2. 分批处理,避免长事务for i in range(0, len(dup_groups), batch_size):batch = dup_groups[i:i + batch_size]with get_db_cursor() as cursor:# 构造 DELETE 语句# 对于 MySQL,可以直接使用 NOT IN 或子查询# 对于 Postgres,需要使用 CTE 或 JOINif config.DB_CONFIG[engine] == mysql:# 利用主键索引加速删除# 注意:这里假设 id 是主键delete_sql = fDELETE FROM {table_name}WHERE id NOT IN (SELECT max_id FROM (SELECT MAX(id) as max_idFROM {table_name}WHERE ({columns_str} = %s, %s, ...)) AS tmp)# 简化逻辑:直接针对每个 group 执行删除# 为了演示通用性,这里使用更稳妥的逐个 Group 处理逻辑for group_vals, max_id, cnt in batch:placeholders = , .join([%s] * len(group_vals))# 删除该组中 ID 不等于 max_id 的记录del_sql = fDELETE FROM {table_name}WHERE ({columns_str}) IN (({placeholders}))AND id != %scursor.execute(del_sql, (*group_vals, max_id))total_deleted += cursor.rowcountelif config.DB_CONFIG[engine] == postgres:# Postgres 语法差异for group_vals, max_id, cnt in batch:placeholders = , .join([%s] * len(group_vals))del_sql = fDELETE FROM {table_name}WHERE ({columns_str}) IN (({placeholders}))AND id != %scursor.execute(del_sql, (*group_vals, max_id))total_deleted += cursor.rowcount# 每批提交一次# MySQL 和 PG 都需要显式提交if config.DB_CONFIG[engine] == mysql:cursor.connection.commit()else:cursor.connection.commit()logger.info(fProcessed batch {i//batch_size + 1}. Deleted {cursor.rowcount} rows in last batch.)logger.info(fCleanup complete. Total deleted: {total_deleted})return total_deleted源码解析关键点:为什么分批? 如果一次删除 10 万条,事务日志(Redo Log/WAL)会激增,可能撑爆磁盘或导致主从延迟。分批提交可以将长事务拆分为短事务,降低锁持有时间。
为什么保留 MAX(id)? 在自增 ID 场景中,ID 越大通常代表数据越新。保留最大 ID 是最常见的业务需求。如果业务需要保留最小 ID,只需将 MAX 改为 MIN。
事务隔离级别:我们在 db_utils 中禁用了 autocommit。这意味着在 commit() 之前,所有 DELETE 操作都在一个事务中。如果中途出错,rollback() 会回滚所有操作,保证数据一致性。运行与测试:复现真实场景
光看代码不够,得跑起来。我们编写 test_data.py 生成脏数据。
# test_data.py
from db_utils import get_db_cursor
import random
import stringdef create_test_table():sql = CREATE TABLE IF NOT EXISTS users (id INT AUTO_INCREMENT PRIMARY KEY,email VARCHAR(255) NOT NULL,name VARCHAR(100),created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP);# Postgres 语法不同,这里仅展示 MySQL 逻辑,PG 需适配 SERIAL 类型with get_db_cursor() as cursor:cursor.execute(DROP TABLE IF EXISTS users)cursor.execute(sql)cursor.connection.commit()def insert_test_data(count=10000):with get_db_cursor() as cursor:for i in range(count):# 故意制造重复:每 10 条中随机选一个重复邮箱base_email = fuser{random.randint(0, 900)}@test.comname = ''.join(random.choices(string.ascii_uppercase, k=5))cursor.execute(INSERT INTO users (email, name) VALUES (%s, %s),(base_email, name))cursor.connection.commit()print(fInserted {count} rows with intentional duplicates.)if __name__ == __main__:create_test_table()insert_test_data()运行 main.py:
# main.py
from duplicate_cleaner import delete_duplicatesif __name__ == __main__:table = users# 假设 email 是业务唯一键,但数据库层面没有唯一约束unique_cols = [email]deleted_count = delete_duplicates(table, unique_cols, batch_size=500)print(fTotal records removed: {deleted_count})测试验证:
运行后,检查 users 表:
SELECT email, COUNT(*) from users GROUP BY email HAVING COUNT(*) 1;结果应为空。同时,抽查某个邮箱,确认保留的是 ID 最大的那条记录。
优化扩展与避坑指南
在 Stack Overflow 上,关于“删除重复数据”的问题高达数千个。总结下来,常见的坑有这三个:索引缺失:如果 email 字段没有索引,GROUP BY 会导致全表扫描,大表直接卡死。对策:在执行清洗前,确保 unique_columns 上有组合索引。
CREATE INDEX idx_email ON users(email);外键约束:如果 id 被其他表引用,直接 DELETE 会报错。对策:在清洗前,检查外键依赖,或使用 ON DELETE CASCADE(慎用),或者先软删除(标记 is_deleted=1),再异步物理删除。并发写入:清洗过程中,如果有新数据插入相同的重复值,会导致刚删完又出现重复。对策:在清洗期间,暂停写入服务,或在应用层加分布式锁。对于高并发场景,建议先 UPDATE 标记,再异步清理,最后 DELETE 标记数据,实现无感清洗。进阶技巧:使用视图辅助
如果重复逻辑非常复杂,可以创建一个临时视图来标识重复项:
CREATE TEMPORARY VIEW dup_users AS
SELECT id, email, ROW_NUMBER() OVER (PARTITION BY email ORDER BY id DESC) as rn
FROM users;DELETE FROM users
WHERE id IN (SELECT id FROM dup_users WHERE rn 1);这种写法在 PostgreSQL 中非常优雅,利用了窗口函数 ROW_NUMBER(),代码可读性远优于嵌套子查询。
小结
删除重复数据保留一条看似简单,实则涉及数据库引擎差异、事务管理、索引优化和并发控制。通过本次实战,我们不仅完成了一个可运行的 Python 工具,更通过源码解析揭示了背后的工程细节。
记住,不要迷信一条 SQL 解决所有问题。在生产环境中,批量处理、事务控制、日志监控缺一不可。配置环境卡半天?按照本文的目录结构和依赖版本,你应该能 5 分钟内跑通。
技术路上没有银弹,只有不断的踩坑与填坑。如果你在项目中遇到了更奇葩的重复数据场景,或者对窗口函数在删除操作中的应用有疑问,还有什么不懂的?评论区留言挨个回。