如何优化从150个CSV文件生成指定结构DataFrame的效率?
优化方案拆解
哇,150个文件花150秒确实有点拖后腿,咱们一步步拆解优化点,应该能把时间压下来不少!
1. 优化CSV读取速度(IO瓶颈)
磁盘IO通常是这类任务的最大瓶颈,先从这里下手:
- 明确指定数据类型:避免pandas自动推断类型的额外开销。比如
site列如果是类别有限的字符串,用'category'类型能大幅节省内存;time列直接用parse_dates解析为datetime,不用后续再转换:df = pd.read_csv(file_path, dtype={'site': 'category'}, parse_dates=['time'], low_memory=False) - 换用更快的读取引擎:pandas支持
engine='pyarrow'(需要先安装pyarrow:pip install pyarrow),读取速度比默认的c引擎快数倍:df = pd.read_csv(file_path, engine='pyarrow', dtype={'site': 'category'}, parse_dates=['time']) - 用Polars替代Pandas:Polars是Rust编写的高性能数据处理库,读取和计算速度远超pandas,语法接近pandas,迁移成本极低:
import polars as pl df = pl.read_csv(file_path, dtypes={'site': pl.Categorical}, parse_dates=['time'])
2. 向量化数据处理(CPU瓶颈)
别用手动循环统计或拆分数据,用内置的矢量化函数,底层都是优化过的C/Rust实现:
- 生成频率字典:用
value_counts()替代手动计数,效率差好几个量级:# Pandas方式 freq_series = df['site'].value_counts() freq_dict = {site: [site, count] for site, count in freq_series.items()} # Polars方式 freq_df = df['site'].value_counts() freq_dict = {row['site']: [row['site'], row['count']] for row in freq_df.to_dicts()} - 拆分会话数据:用numpy的
reshape一次性完成10条一组的拆分,替代逐行循环:# 先按时间排序(会话必须是时间顺序) df = df.sort_values('time') # 计算完整会话数(丢弃不足10条的尾部) n_sessions = len(df) // 10 if n_sessions == 0: continue # 一次性拆分成n_sessions行×10列的数组 session_array = df['site'].values[:n_sessions*10].reshape(n_sessions, 10) # 转为DataFrame session_df = pd.DataFrame(session_array, columns=[f'site{i+1}' for i in range(10)]) session_df['user_id'] = user_id
3. 并行处理多个文件
150个文件是完全独立的,正好用多进程并行处理,榨干多核CPU的性能:
- 用
concurrent.futures.ProcessPoolExecutor:
注:如果是纯IO密集型场景,也可以用from concurrent.futures import ProcessPoolExecutor import pandas as pd import glob import os def process_single_file(file_path): # 解析user_id,用os模块更适配不同操作系统 user_id = os.path.splitext(os.path.basename(file_path))[0] # 读取文件 df = pd.read_csv(file_path, engine='pyarrow', dtype={'site': 'category'}, parse_dates=['time']) # 排序并拆分会话 df = df.sort_values('time') n_sessions = len(df) // 10 if n_sessions == 0: return None session_array = df['site'].values[:n_sessions*10].reshape(n_sessions, 10) session_df = pd.DataFrame(session_array, columns=[f'site{i+1}' for i in range(10)]) session_df['user_id'] = user_id # 生成频率字典 freq_dict = {site: [site, count] for site, count in df['site'].value_counts().items()} return session_df, freq_dict # 获取所有文件路径 file_paths = glob.glob('../user*.csv') # 并行处理,默认用CPU核心数 with ProcessPoolExecutor() as executor: results = list(executor.map(process_single_file, file_paths)) # 合并最终结果 all_session_dfs = [res[0] for res in results if res[0] is not None] final_session_df = pd.concat(all_session_dfs, ignore_index=True) all_freq_dicts = {} for res in results: if res[1] is not None: all_freq_dicts.update(res[1])ThreadPoolExecutor,但多进程更适合CPU+IO混合的场景。
4. 砍掉不必要的操作
- 全局统计频率字典:如果不需要每个文件单独的频率字典,可以合并所有用户的site数据后统一统计,减少多次统计的开销:
# 先合并所有会话的site列 all_sites = final_session_df[[f'site{i+1}' for i in range(10)]].stack() global_freq_dict = {site: [site, count] for site, count in all_sites.value_counts().items()} - 复用对象:提前定义列名列表,别每次循环都生成:
SITE_COLUMNS = [f'site{i+1}' for i in range(10)] # 后面直接用SITE_COLUMNS session_df = pd.DataFrame(session_array, columns=SITE_COLUMNS)
5. 硬件/文件系统优化
- 如果文件在机械硬盘上,复制到**固态硬盘(SSD)**上,磁盘寻道时间会大幅降低,读取速度能翻好几倍。
- 如果后续经常要处理这些文件,可以把所有CSV合并成一个大文件(按user_id分组),减少多次打开/关闭文件的开销。
按照这些优化点,把处理时间从150秒降到几十秒甚至10秒以内完全没问题!
内容的提问来源于stack exchange,提问作者Денис Милованов
相关产品推荐
相关产品推荐

