如何在不丢失过期记录的前提下增量合并每日CSV Delta更新至主文件?
背景与数据结构
- 数据源:每日生成的CSV文件,记录预订状态变更
- 核心字段:
Booking_ID、Previous_Status、New_Status、Timestamp(格式YYYY-MM-DD HH:MM:SS) - 边缘场景:
- 单个
Booking_ID单日可有多条状态变更记录 - 同一
Booking_ID的不同变更记录可能有完全相同的Timestamp
- 单个
面临挑战
源系统会自动清理过期预订记录,因此不能每日全量覆盖或重新加载,只能处理每日的Delta文件(仅含新预订或状态变更记录)并合并到主数据集。脚本需每日运行,仅读取目录中新添加的CSV文件以保证性能。
已解决问题
通过基于系统 ingestion 顺序的自定义序列索引,结合Timestamp实现了同一Booking_ID下重复时间戳记录的正确排序。
当前问题
在累积合并过程中,如何高效对比Delta文件与主数据集:
- 追加全新的预订记录
- 按正确顺序插入新的状态变更记录
- 避免重复处理已存在的变更记录
寻求成熟的模式、算法或Python/Pandas/SQL等工具的特性来实现这类增量式 ledger-style 合并,以及相关架构建议。
Python/Pandas 实现方案
构建唯一标识键
利用已有的Timestamp+自定义序列索引,结合Booking_ID生成唯一主键(例如命名为Change_Key,格式为Booking_ID|Timestamp|Sequence_Index)。主数据集和Delta文件都生成该键,用于快速判断记录是否已存在。增量合并流程
- 读取当日新增的Delta文件,生成
Change_Key - 读取主数据集的
Change_Key集合(无需全量加载主数据,可单独维护键的索引文件或用轻量数据库存储) - 筛选Delta文件中
Change_Key不在主数据集键集合内的记录 - 将筛选后的新记录追加到主数据集
- 按
Booking_ID、Timestamp、Sequence_Index排序,保证每个预订的状态流顺序
代码示例:
import pandas as pd import os # 配置路径 MAIN_DATA_PATH = "main_booking_data.csv" DELTA_DIR = "daily_deltas/" # 筛选当日新增Delta文件(可按命名规则/创建时间判断) delta_files = [f for f in os.listdir(DELTA_DIR) if f.startswith("delta_2024")] delta_df = pd.concat([pd.read_csv(os.path.join(DELTA_DIR, f)) for f in delta_files]) # 生成唯一标识键 delta_df["Change_Key"] = delta_df.apply( lambda x: f"{x['Booking_ID']}|{x['Timestamp']}|{x['Sequence_Index']}", axis=1 ) # 合并到主数据集 if os.path.exists(MAIN_DATA_PATH): main_df = pd.read_csv(MAIN_DATA_PATH) existing_keys = set(main_df["Change_Key"]) new_records = delta_df[~delta_df["Change_Key"].isin(existing_keys)] combined_df = pd.concat([main_df, new_records], ignore_index=True) else: # 首次运行直接用Delta数据初始化主数据集 combined_df = delta_df # 按预订和时间顺序排序 combined_df = combined_df.sort_values( by=["Booking_ID", "Timestamp", "Sequence_Index"], ignore_index=True ) combined_df.to_csv(MAIN_DATA_PATH, index=False)- 读取当日新增的Delta文件,生成
性能优化
- 单独维护
change_keys.txt或用SQLite存储键集合,避免每次全量读取主数据集 - 超大数据集场景下,用Dask替代Pandas做并行化处理
- 单独维护
SQL 实现方案
用数据库存储主数据集能大幅提升合并效率:
创建主表并建立唯一索引
CREATE TABLE booking_status_changes ( Booking_ID VARCHAR(50), Previous_Status VARCHAR(20), New_Status VARCHAR(20), Timestamp DATETIME, Sequence_Index INT, Change_Key VARCHAR(100) PRIMARY KEY -- 由Booking_ID+Timestamp+Sequence_Index拼接生成 );增量合并(以MySQL为例)
将Delta文件导入临时表后,利用主键冲突跳过重复记录:-- 先将Delta数据导入临时表temp_booking_changes INSERT INTO booking_status_changes (Booking_ID, Previous_Status, New_Status, Timestamp, Sequence_Index, Change_Key) SELECT Booking_ID, Previous_Status, New_Status, Timestamp, Sequence_Index, Change_Key FROM temp_booking_changes ON DUPLICATE KEY UPDATE Booking_ID = Booking_ID; -- 主键冲突时不执行更新,直接跳过重复记录查询完整状态流
直接按预订ID和时间排序即可获取完整状态变更历史:SELECT * FROM booking_status_changes WHERE Booking_ID = 'BK12345' ORDER BY Timestamp, Sequence_Index;
成熟架构模式建议
事件溯源(Event Sourcing)
将所有状态变更作为不可变事件存储,主数据集即为事件的有序集合。每个事件分配唯一标识(全局或预订维度),只需保证事件不重复写入,天然适配增量合并场景。CDC风格批次管理
给每个Delta文件打唯一批次标识,记录已处理的批次列表,避免重复处理整个文件;结合单条记录的唯一键去重,实现“恰好一次”的处理语义。分层存储
- 热存储:用PostgreSQL、SQLite等数据库存储最近的变更记录,方便快速查询与合并
- 冷存储:将历史归档记录存为Parquet/CSV文件,降低热存储的资源压力
内容的提问来源于stack exchange,提问作者N M Dako

