如何高效拼接数千个Pandas DataFrame并批量处理CSV源文件
CSV合并处理优化方案
原实现核心问题
- 全量文件重复读写:每处理100个文件就读取已存储的全量feather数据再拼接写入,随数据量增长读写开销呈指数级上升
- 内存无意义占用:全程维护
concat_df变量累加所有已处理数据,未及时释放中间缓存,内存占用随处理进度持续上涨 - 存在代码bug:循环读取CSV时错误使用未定义的
file_name变量,而非迭代变量csv_name
优化后实现
总数据量为7000文件 * 600行 = 420万行,仅占用百MB级内存,常规设备完全可以直接全量处理,以下是最优实现:
import os import pandas as pd from datetime import datetime, timedelta # 自定义配置 CSV_DIR = "/data/csvs" OUTPUT_PATH = "./merged_result.feather" # 每行时间增量为0.5秒 ROW_TIME_DELTA = timedelta(seconds=0.5) # 预筛选所有CSV文件 csv_file_list = [os.path.join(CSV_DIR, f) for f in os.listdir(CSV_DIR) if f.endswith(".csv")] def add_timestamp_column(df, file_name): # 解析文件名中的日期:格式为DDMMYY date_part = os.path.splitext(file_name)[0] base_datetime = datetime.strptime(date_part, "%d%m%y") # 向量化生成时间序列,无循环性能更高 df["timestamp"] = base_datetime + pd.to_timedelta(df.index * ROW_TIME_DELTA) # 若需要字符串格式的时分秒微秒,可打开下一行注释 # df["timestamp"] = df["timestamp"].dt.strftime("%H:%M:%S.%f") return df # 用生成器迭代处理文件,内存仅保留当前单文件数据 df_iter = (add_timestamp_column(pd.read_csv(file_path), os.path.basename(file_path)) for file_path in csv_file_list) # 一次性合并所有数据 full_df = pd.concat(df_iter, ignore_index=True) # 输出最终结果 full_df.to_feather(OUTPUT_PATH)
极端低内存场景适配
如果设备内存小于1GB,可采用分片写临时文件再合并的方案,全程内存占用不超过100MB:
import os import glob import pandas as pd from datetime import datetime, timedelta # 自定义配置 CSV_DIR = "/data/csvs" TEMP_DIR = "./temp_batch" OUTPUT_PATH = "./merged_result.feather" ROW_TIME_DELTA = timedelta(seconds=0.5) BATCH_SIZE = 100 # 每100个文件写一次临时分片 os.makedirs(TEMP_DIR, exist_ok=True) csv_file_list = [os.path.join(CSV_DIR, f) for f in os.listdir(CSV_DIR) if f.endswith(".csv")] def add_timestamp_column(df, file_name): date_part = os.path.splitext(file_name)[0] base_datetime = datetime.strptime(date_part, "%d%m%y") df["timestamp"] = base_datetime + pd.to_timedelta(df.index * ROW_TIME_DELTA) return df # 分批次写入临时文件 for batch_idx in range(0, len(csv_file_list), BATCH_SIZE): batch_files = csv_file_list[batch_idx: batch_idx+BATCH_SIZE] batch_df_list = [] for path in batch_files: df = pd.read_csv(path) df = add_timestamp_column(df, os.path.basename(path)) batch_df_list.append(df) batch_full = pd.concat(batch_df_list, ignore_index=True) batch_full.to_feather(os.path.join(TEMP_DIR, f"batch_{batch_idx}.feather")) del batch_df_list, batch_full # 合并所有临时分片 temp_file_list = glob.glob(os.path.join(TEMP_DIR, "*.feather")) full_df = pd.concat((pd.read_feather(p) for p in temp_file_list), ignore_index=True) full_df.to_feather(OUTPUT_PATH) # 清理临时文件 for p in temp_file_list: os.remove(p) os.rmdir(TEMP_DIR)
内容的提问来源于stack exchange,提问作者georgeamccarthy
相关产品推荐
相关产品推荐

