基于每周CSV Dump,如何用Pandas/SQL识别PostgreSQL中被删除的记录?
识别CSV增量更新中被删除的记录(Pandas/SQL实现方案)
问题背景
我有一个Django应用,负责处理包含1.5年数据的每周CSV Dump,并将数据存储到PostgreSQL数据库中。数据包含主键列、uploadDate列及其他字段。当源端删除某条记录时,新的CSV Dump中将不再包含该记录,但它仍存在于我的数据库内。需要通过Pandas或SQL,按uploadDate升序排序数据,对比相邻的旧文件与新文件记录,找出旧文件存在但新文件缺失的被删除记录。
现有Pandas代码评估
你提供的代码尝试通过分组主键并检查首次/最后出现日期之间的记录存在性来识别删除,但存在以下问题:
- 逻辑不准确:若某个主键临时消失后又重新出现(如周1存在、周2消失、周3重现),会被误判为已删除。
- 效率低下:循环遍历每个主键分组,在数据量较大时(如百万级主键)运行速度极慢,未利用Pandas向量化操作的优势。
- 信息不完整:仅返回主键和最后出现日期,无法明确删除发生的具体时间点。
import pandas as pd def find_deleted_records(df: pd.DataFrame, primary_key: str, upload_date: str) -> pd.DataFrame: """ Find records that were deleted from a database table based on a weekly CSV dump. Args: df (pd.DataFrame): The dataframe containing the table data. primary_key (str): The name of the primary key column. upload_date (str): The name of the upload date column. Returns: pd.DataFrame: A dataframe containing the deleted records. """ # Sort the data by Peildatum in ascending order data = df.sort_values(by=upload_date) # Group the data by Uniek kenmerk and get the first and last Peildatum for each group grouped_data = data.groupby(primary_key).agg({upload_date: ['first', 'last']}) # Identify the records that were deleted deleted_records = [] for index, row in grouped_data.iterrows(): first_Peildatum = row[upload_date]['first'] last_Peildatum = row[upload_date]['last'] mask = (data[primary_key] == index) & (data[upload_date] > first_Peildatum) & (data[upload_date] < last_Peildatum) if len(data.loc[mask]) == 0: deleted_records.append({primary_key: index, upload_date: last_Peildatum}) # Return the deleted records as a dataframe return pd.DataFrame(deleted_records)
更优Pandas实现方案
核心思路
- 提取所有唯一的
uploadDate并按升序排序。 - 遍历每一对相邻的日期,对比前后两个日期的主键集合。
- 找出前一个日期存在、后一个日期缺失的主键(即本周被删除的记录),并避免重复记录同一主键。
代码实现
import pandas as pd def find_deleted_records(df: pd.DataFrame, primary_key: str, upload_date: str) -> pd.DataFrame: # 按uploadDate升序获取唯一日期列表 sorted_dates = sorted(df[upload_date].unique()) deleted_records = [] deleted_keys = set() # 存储已识别的删除主键,避免重复 for i in range(1, len(sorted_dates)): prev_date = sorted_dates[i-1] curr_date = sorted_dates[i] # 获取前后日期对应的主键集合 prev_keys = set(df[df[upload_date] == prev_date][primary_key]) curr_keys = set(df[df[upload_date] == curr_date][primary_key]) # 筛选出新增的被删除主键(前有后无且未被记录过) newly_deleted = prev_keys - curr_keys - deleted_keys deleted_keys.update(newly_deleted) # 整理结果数据 for key in newly_deleted: deleted_records.append({ primary_key: key, "last_seen_date": prev_date, "deleted_detected_date": curr_date }) return pd.DataFrame(deleted_records)
方案优势
- 逻辑精准:仅对比相邻周的记录,有效避免误判临时消失后重现的主键。
- 效率更高:利用集合运算替代循环遍历,处理大数据量时性能提升明显。
- 信息完整:返回主键、最后出现日期和删除检测日期,便于后续数据库清理或跟踪。
PostgreSQL SQL实现方案
方案1:找出最终被删除的记录(适合数据库清理)
此SQL会找出在最新CSV中不存在的所有记录,即源端已删除且未恢复的记录:
WITH all_upload_dates AS ( SELECT DISTINCT uploadDate FROM your_table ORDER BY uploadDate ), latest_upload_date AS ( SELECT MAX(uploadDate) AS max_date FROM all_upload_dates ), key_last_seen AS ( SELECT primary_key, MAX(uploadDate) AS last_seen_date FROM your_table GROUP BY primary_key ) SELECT kls.primary_key, kls.last_seen_date, lud.max_date AS deleted_detected_date FROM key_last_seen kls CROSS JOIN latest_upload_date lud WHERE NOT EXISTS ( SELECT 1 FROM your_table t WHERE t.primary_key = kls.primary_key AND t.uploadDate = lud.max_date );
方案2:按周跟踪删除记录(与Pandas方案逻辑一致)
此SQL会找出每一周新增的被删除记录,并仅保留每个主键首次被检测到删除的记录:
WITH date_pairs AS ( SELECT uploadDate AS curr_date, LAG(uploadDate) OVER (ORDER BY uploadDate) AS prev_date FROM (SELECT DISTINCT uploadDate FROM your_table) AS dates WHERE LAG(uploadDate) OVER (ORDER BY uploadDate) IS NOT NULL ), prev_period_keys AS ( SELECT dp.prev_date, t.primary_key FROM date_pairs dp JOIN your_table t ON t.uploadDate = dp.prev_date ), curr_period_keys AS ( SELECT dp.curr_date, t.primary_key FROM date_pairs dp JOIN your_table t ON t.uploadDate = dp.curr_date ) SELECT pk.prev_date AS last_seen_date, pk.curr_date AS deleted_detected_date, pk.primary_key FROM ( SELECT dp.prev_date, dp.curr_date, ppk.primary_key FROM date_pairs dp JOIN prev_period_keys ppk ON ppk.prev_date = dp.prev_date WHERE NOT EXISTS ( SELECT 1 FROM curr_period_keys cpk WHERE cpk.curr_date = dp.curr_date AND cpk.primary_key = ppk.primary_key ) ) AS pk QUALIFY ROW_NUMBER() OVER (PARTITION BY pk.primary_key ORDER BY pk.curr_date) = 1;
内容的提问来源于stack exchange,提问作者Abashirsoldom
相关产品推荐
相关产品推荐

