Jupyter Notebook批量读取大文件并分批分析方案咨询
分批处理大文件生成特征矩阵的解决方案
完全可以通过分批读取、分析并保存中间结果的方式解决内存溢出问题,以下是结合你的代码框架的具体实现方案:
核心思路
- 先整理所有目标文件的路径,按固定批次大小(比如10个)拆分
- 将你的分析逻辑封装为可重复调用的函数,避免重复编写代码
- 循环处理每一批文件,生成并保存该批次的特征矩阵
- 最后合并所有批次的结果得到完整特征矩阵
实现代码
import pathlib import pandas as pd # 定义每批处理的文件数量 BATCH_SIZE = 10 # ---------------------- # 1. 整理所有待处理的文件路径 # ---------------------- # 处理文件夹A的路径与文件列表 paths_A = [ pathlib.Path("/directory/for/path/1A/"), pathlib.Path("/directory/for/path/2A/"), pathlib.Path("/directory/for/path/3A/") ] # 过滤出每个目录下的CSV文件(避免读取非文件或非CSV格式内容) file_groups_A = [] for p in paths_A: csv_files = sorted(list(p.glob("*.csv"))) # 排序保证文件匹配一致性 file_groups_A.append(csv_files) # 处理文件夹B的路径与文件列表 paths_B = [ pathlib.Path("/directory/for/path/1B/"), pathlib.Path("/directory/for/path/2B/"), pathlib.Path("/directory/for/path/3B/") ] file_groups_B = [] for p in paths_B: csv_files = sorted(list(p.glob("*.csv"))) file_groups_B.append(csv_files) # 确保所有文件组的文件数量一致(避免批次处理时出现不匹配) max_file_count = min(len(files) for files in file_groups_A + file_groups_B) file_groups_A = [files[:max_file_count] for files in file_groups_A] file_groups_B = [files[:max_file_count] for files in file_groups_B] # ---------------------- # 2. 封装分析逻辑为函数 # ---------------------- def process_single_batch(batch_files_A, batch_files_B): # 读取当前批次的文件到对应列表 list_1_A, list_2_A, list_3_A = [], [], [] for f1, f2, f3 in zip(*batch_files_A): # 替换为你实际的文件读取与数据提取逻辑 df1 = pd.read_csv(f1) list_1_A.append(df1) df2 = pd.read_csv(f2) list_2_A.append(df2) df3 = pd.read_csv(f3) list_3_A.append(df3) list_1_B, list_2_B, list_3_B = [], [], [] for f1, f2, f3 in zip(*batch_files_B): # 替换为你实际的文件读取与数据提取逻辑 df1 = pd.read_csv(f1) list_1_B.append(df1) df2 = pd.read_csv(f2) list_2_B.append(df2) df3 = pd.read_csv(f3) list_3_B.append(df3) # ---------------------- # 插入你原来100+单元格的分析代码 # 以list_1_A、list_2_A等为输入,最终生成feature_matrix # 示例分析逻辑(替换为你的实际代码): feature_matrix = pd.concat([ pd.concat(list_1_A).agg(["mean", "sum"]), pd.concat(list_2_A).agg(["min", "max"]), pd.concat(list_3_A).agg(["median", "std"]), pd.concat(list_1_B).agg(["mean", "sum"]), pd.concat(list_2_B).agg(["min", "max"]), pd.concat(list_3_B).agg(["median", "std"]) ], axis=1).T # ---------------------- return feature_matrix # ---------------------- # 3. 分批处理并保存中间结果 # ---------------------- total_batches = (max_file_count + BATCH_SIZE - 1) // BATCH_SIZE # 向上取整计算总批次 for batch_idx in range(total_batches): start_idx = batch_idx * BATCH_SIZE end_idx = start_idx + BATCH_SIZE # 获取当前批次的文件子集 current_batch_A = [files[start_idx:end_idx] for files in file_groups_A] current_batch_B = [files[start_idx:end_idx] for files in file_groups_B] # 跳过空批次(最后一批可能不足BATCH_SIZE) if not current_batch_A[0]: continue # 处理当前批次 batch_feature_matrix = process_single_batch(current_batch_A, current_batch_B) # 保存批次结果(推荐用Parquet格式,比CSV更节省空间且读写更快) batch_feature_matrix.to_csv(f"batch_feature_matrix_{batch_idx+1}.csv", index=False) # 可选:batch_feature_matrix.to_parquet(f"batch_feature_matrix_{batch_idx+1}.parquet") print(f"批次 {batch_idx+1} 处理完成,已保存中间结果") # ---------------------- # 4. 合并所有批次结果 # ---------------------- batch_dataframes = [] for batch_idx in range(total_batches): try: # 读取批次文件 df = pd.read_csv(f"batch_feature_matrix_{batch_idx+1}.csv") # 可选:df = pd.read_parquet(f"batch_feature_matrix_{batch_idx+1}.parquet") batch_dataframes.append(df) except FileNotFoundError: continue # 合并为最终特征矩阵 final_feature_matrix = pd.concat(batch_dataframes, ignore_index=True) final_feature_matrix.to_csv("final_feature_matrix.csv", index=False) # 可选:final_feature_matrix.to_parquet("final_feature_matrix.parquet") print("所有批次合并完成,最终特征矩阵已保存")
关键注意事项
- 文件匹配一致性:使用
sorted()对文件列表排序,确保每个文件夹下的文件能一一对应(比如按文件名排序) - 内存优化:读取CSV时可通过
pd.read_csv的dtype参数指定数据类型(如将字符串列设为category,数值列设为float32),或用chunksize拆分单个超大CSV文件 - 错误处理:可在文件读取和分析逻辑中加入
try-except块,避免单个文件出错导致整个批次失败 - 中间文件管理:合并完成后可删除批次文件,节省磁盘空间
内容的提问来源于stack exchange,提问作者S C
相关产品推荐
相关产品推荐

