基于Python MySQL Connector的ETL事实表去重方案咨询
Python MySQL 销售事实表去重方案评估与优化建议
现有方案评估
方案1:SQL查询已存在记录后本地去重
- 优势:逻辑直接,基于数据库源数据去重,不会出现数据不一致问题
- 劣势:针对200万行的大表,每次全量拉取
InvoiceDate和InvoiceNumber会占用大量数据库带宽和本地内存,随着数据量增长,这个步骤会成为ETL流程的性能瓶颈,自动化运行时容易出现超时或资源不足问题
方案2:本地CSV校验文件
- 优势:本地文件IO速度快,能减少数据库查询压力
- 劣势:健壮性极差,一旦ETL中断、文件损坏或遗漏数据,必然导致重复插入;可扩展性差,校验窗口扩大后CSV文件会急剧膨胀,多节点部署时无法共享校验文件,极易出现数据不一致
更优方案建议
1. 数据库唯一约束 + 插入时忽略重复(推荐)
这是最可靠且高效的方案,将去重逻辑交给数据库处理:
- 第一步:给
sales表添加InvoiceDate和InvoiceNumber的组合唯一约束:ALTER TABLE sales ADD UNIQUE KEY idx_invoice_unique (InvoiceDate, InvoiceNumber); - 第二步:修改Python插入逻辑,使用
INSERT IGNORE实现批量插入时自动忽略重复:
示例代码(使用MySQL原生连接批量插入):
这种方式无需拉取大量数据到本地,数据库层面直接拦截重复数据,性能和健壮性拉满,完全适配自动化ETL流程。import pandas as pd from sqlalchemy import create_engine engine = create_engine('mysql+mysqlconnector://[用户名]:[密码]@[主机地址]/[数据库名]') def batch_insert_ignore(df, table_name): conn = engine.raw_connection() cursor = conn.cursor() # 构造插入语句 column_str = ','.join(df.columns) placeholder_str = ','.join(['%s'] * len(df.columns)) insert_sql = f"INSERT IGNORE INTO {table_name} ({column_str}) VALUES ({placeholder_str})" # 批量执行插入 cursor.executemany(insert_sql, df.values.tolist()) conn.commit() cursor.close() conn.close() # 调用方法插入清洗后的销售数据 batch_insert_ignore(sales, 'sales')
2. 增量加载 + 时间范围过滤
如果原始数据源支持按时间增量获取(比如每日仅拉取当天数据):
- 调整ETL第一步,只读取最近N天的原始数据,而非全量数据
- 查询数据库中对应时间范围的已存在
InvoiceDate和InvoiceNumber,过滤本地DataFrame中的重复记录后再插入 - 这种方案大幅减少了需要处理的数据量,兼顾准确性和性能,适合有明确时间维度的销售数据场景
3. 临时表中转批量去重
针对超大批量数据的场景,利用数据库的批量处理能力:
- 将清洗后的全量数据插入MySQL临时表
temp_sales(临时表会话结束后自动销毁) - 执行SQL语句完成去重插入:
INSERT INTO sales SELECT t.* FROM temp_sales t LEFT JOIN sales s ON t.InvoiceDate = s.InvoiceDate AND t.InvoiceNumber = s.InvoiceNumber WHERE s.InvoiceNumber IS NULL; - 数据库的join操作性能远优于Python内存中的数据处理,能快速完成百万级数据的去重插入
内容的提问来源于stack exchange,提问作者ProfessorE
相关产品推荐
相关产品推荐

