You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何优化从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:
    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])
    
    注:如果是纯IO密集型场景,也可以用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,提问作者Денис Милованов

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 08:24:09