每日10分钟间隔Dump的DataFrame增量合并高效实现方案问询
高效处理每日增量DataFrame Dump文件的优化方案
针对你这种拥有5000个按10分钟间隔生成的DataFrame Dump文件的场景,直接全量读取后去重确实会浪费大量计算资源——毕竟同一天的相邻文件只有约2行差异,全量处理重复内容完全没必要。下面给你一个按日期分组的高效实现方案,精准处理每日增量,同时利用Dask的并行能力提升效率:
核心思路
- 按日期分组文件:从文件名中提取日期,把同一天的文件归为一组,避免跨日期的无效重复处理
- 每日内部去重:对当天所有文件的合并数据去重,得到当天包含所有新增行的完整数据集(等价于从首个文件开始依次追加后续新增行的结果)
- 跨日期合并:把所有日期的处理结果合并成最终的完整DataFrame
这个思路把重复数据的处理范围缩小到单日内部,比全量去重的计算量要小得多,尤其适合你这种每日增量极少的场景。
具体实现代码
import pandas as pd import dask.dataframe as dd from glob import glob from collections import defaultdict # 1. 获取所有文件路径并按日期分组 file_paths = glob("path/to/your/files/*.csv") # 替换成你的文件路径匹配规则 date_groups = defaultdict(list) for path in file_paths: # 从文件名中提取日期(根据你的文件名格式调整提取逻辑) # 示例文件名:2019-08-28 06:00:13 SCHOOL_20190828... date_str = path.split()[0] # 从开头的时间戳提取日期 # 如果文件名格式特殊,也可以用SCHOOL_后的8位日期提取: # date_part = path.split("SCHOOL_")[1][:8] # date_str = f"{date_part[:4]}-{date_part[4:6]}-{date_part[6:8]}" date_groups[date_str].append(path) # 2. 按日期处理每组文件 daily_results = [] for date, files in date_groups.items(): print(f"Processing date: {date}") # 读取当天所有文件为Dask DataFrame ddf_daily = dd.read_csv( files, storage_options=storage_opts, # 替换成你的存储配置(云存储/本地) assume_missing=True ) # 对当天数据去重——得到当天所有唯一行(等价于从首个文件追加所有新增行) ddf_daily_unique = ddf_daily.drop_duplicates() # 计算得到当天的结果DataFrame df_daily = ddf_daily_unique.compute() daily_results.append(df_daily) # 3. 合并所有日期的结果 final_df = pd.concat(daily_results, ignore_index=True) # 如果不同日期可能存在重复行,可追加跨日期去重: # final_df = final_df.drop_duplicates(ignore_index=True)
为什么这个方案更高效?
- 减少重复计算:每天内部先去重,避免了把同一天的重复数据带到跨日期处理中,相比全量去重,计算量至少能减少一个数量级(毕竟每天有144个文件,大部分内容重复)
- Dask并行优势:Dask会自动并行处理当天的多个文件,利用多核CPU提升读取和去重的速度
- 内存友好:单日数据的内存占用远小于全量数据,避免了内存溢出的风险
如果你使用云存储(比如S3、OSS),可以用fsspec代替glob来获取文件路径,逻辑完全一致。
内容的提问来源于stack exchange,提问作者theSekyi
相关产品推荐
相关产品推荐

