并行读取CSV并合并DataFrame的代码问题及优化咨询
嘿,咱们一步步来解决你的并行CSV处理问题吧!我会先拆解你当前代码里的问题,再给出优化方案,让这个批量处理任务更顺畅高效。
批量CSV并行读取与合并的问题解决
原代码的核心问题拆解
- 回调函数无法区分任务来源:你的
concatDf没有接收回调返回的结果,还硬编码要合并df_1和df_2,完全没法适配动态的多文件任务结果。 - 无法判断所有任务完成:原代码没有等待所有线程任务结束,也没有统一收集所有返回的DataFrame,根本没法合并50个文件的结果。
- 实现方式效率低下:
ThreadPool本身适合IO密集型任务,但你用apply_sync逐个提交任务的方式太繁琐,而且开40个线程反而会增加线程切换开销,IO密集型任务8-16个线程足够。
问题1&2的解决方案:用批量任务API统一收集结果
其实完全不需要用回调函数来合并,更简单的方式是用pool.map(或现代的ThreadPoolExecutor.map)批量提交任务,它会自动等待所有任务完成,返回所有DataFrame的列表,完美解决“区分任务”和“判断完成”的问题。
修改后的基础实现
import pandas as pd from concurrent.futures import ThreadPoolExecutor from time import time import glob class DataProvider: # 注意Python类名要大驼峰,原代码的小写Class是语法错误 def __init__(self, num_workers=10): self.num_workers = num_workers self.final_df = pd.DataFrame() def read_single_csv(self, filename): # 这里可以提前加读取优化,比如指定数据类型、跳过不需要的列 return pd.read_csv(filename, low_memory=False) def merge_all_csv(self, csv_source): start_time = time() # 获取所有CSV文件路径:可以是glob匹配,也可以是手动传入的文件列表 if isinstance(csv_source, str): csv_files = glob.glob(csv_source) else: csv_files = csv_source # 批量提交任务,map会自动等待所有任务完成并返回结果列表 with ThreadPoolExecutor(max_workers=self.num_workers) as executor: all_dfs = list(executor.map(self.read_single_csv, csv_files)) # 合并所有DataFrame,并提取第2列及以后的内容 self.final_df = pd.concat(all_dfs, ignore_index=True).iloc[:, 1:] total_time = time() - start_time print(f"全部处理完成,总耗时:{total_time:.2f}秒") return self.final_df # 使用示例 if __name__ == "__main__": provider = DataProvider(num_workers=10) # 替换成你的CSV文件路径,比如"*.csv"或者具体的文件列表 merged_df = provider.merge_all_csv("path/to/your/csvs/*.csv")
问题3:更高效的进阶优化方案
针对400MB级别的大CSV,还有几个关键优化点能大幅提升速度和内存利用率:
1. 优化CSV读取参数
在pd.read_csv中添加这些参数,减少内存占用和读取时间:
dtype:提前指定列的数据类型,比如把字符串列设为'category',数值列设为np.float32/np.int32,避免pandas自动推断的开销。usecols:直接跳过不需要的列(比如你最后要iloc[:,1:],读取时就可以跳过第一列),减少内存占用。low_memory=False:避免大文件读取时的类型推断警告,同时加快读取速度。
示例:
def read_single_csv(self, filename): return pd.read_csv( filename, usecols=lambda col: col != 0, # 跳过索引为0的第一列 dtype={"user_id": np.int32, "category": "category"}, low_memory=False )
2. 分阶段合并降低内存峰值
如果50个文件合并后内存压力大,可以分批次合并(比如每10个文件合并一次),再合并中间结果,减少内存峰值:
def merge_all_csv(self, csv_files): start_time = time() batch_size = 10 batches = [csv_files[i:i+batch_size] for i in range(0, len(csv_files), batch_size)] intermediate_dfs = [] with ThreadPoolExecutor(max_workers=self.num_workers) as executor: for batch in batches: batch_dfs = list(executor.map(self.read_single_csv, batch)) intermediate_dfs.append(pd.concat(batch_dfs, ignore_index=True)) self.final_df = pd.concat(intermediate_dfs, ignore_index=True).iloc[:, 1:] print(f"处理完成,总耗时:{time()-start_time:.2f}秒") return self.final_df
3. 超大数据场景:用Dask替代Pandas
如果合并后的DataFrame超出内存,可以用Dask DataFrame,它能并行处理并延迟加载,不需要一次性把所有数据放进内存:
import dask.dataframe as dd def merge_with_dask(csv_pattern): start_time = time() # Dask自动并行读取所有CSV ddf = dd.read_csv(csv_pattern) # 提取第2列及以后的内容 ddf = ddf.iloc[:, 1:] # 转换为Pandas DataFrame(内存足够时),或直接用Dask做后续处理 final_df = ddf.compute() print(f"处理完成,总耗时:{time()-start_time:.2f}秒") return final_df
Dask会自动管理并行任务,不需要手动维护线程池,非常适合超大数据场景。
总结
- 放弃回调函数,用
ThreadPoolExecutor.map批量收集任务结果,天然解决任务区分和完成判断的问题。 - 优化
pd.read_csv的参数,减少内存占用和读取时间。 - 内存压力大时用分阶段合并,超大数据场景直接用Dask。
内容的提问来源于stack exchange,提问作者Invictus
相关产品推荐
相关产品推荐

