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

多次合并DataFrame是否高效?——Python分批处理CSV场景咨询

处理百万行CSV并合并为单个DataFrame生成Parquet文件的优化方案

你提到的分批处理百万级CSV的思路非常务实,完美避开了一次性加载全量数据导致的内存爆炸问题,同时也能满足后续生成单个Parquet文件的需求。下面我会给出完整的实现代码,再补充几个关键优化点,帮你在Hadoop环境下更顺畅地完成任务。

核心实现代码

这里我基于你的思路完善了逻辑,重点优化了DataFrame的合并效率,避免频繁内存复制:

import pandas as pd

class CSVProcessor:
    def __init__(self):
        self.df_list = []  # 用列表暂存小DataFrame,比反复concat更省内存

    def process_batch(self, batch):
        """这里替换成你的实际批处理逻辑:清洗、转换、特征工程等"""
        # 举个例子:筛选有效数据+新增计算列
        processed_batch = batch[batch['amount'] > 0].copy()
        processed_batch['double_amount'] = processed_batch['amount'] * 2
        return processed_batch

    def parallel_process(self, csv_path, chunksize=10000):
        # 利用pandas内置的chunksize分批读取CSV,无需手动拆分
        for chunk in pd.read_csv(csv_path, chunksize=chunksize):
            # 处理当前批次
            processed_chunk = self.process_batch(chunk)
            # 将处理后的小DataFrame存入列表
            self.df_list.append(processed_chunk)
        
        # 最后一次性合并所有小DataFrame,生成最终大DataFrame
        self.df = pd.concat(self.df_list, ignore_index=True)
        # 清空列表释放冗余内存
        self.df_list = None

    def save_to_hadoop_parquet(self, hdfs_path):
        """将合并后的DataFrame写入HDFS的单个Parquet文件"""
        # 推荐用pyarrow引擎,对HDFS支持更好,需提前安装pyarrow库
        self.df.to_parquet(
            hdfs_path,
            engine='pyarrow',
            compression='snappy',  # 压缩比和速度平衡的最优选择
            storage_options={'host': 'your-namenode-ip', 'port': 9000}  # 按需配置HDFS参数
        )

关键优化说明

  • 合并策略优化:不要每次都把新的小DataFrame拼接到self.df里——每次pd.concat([self.df, new_chunk])都会生成全新的DataFrame,重复复制数据会导致内存开销激增。用列表暂存所有小DataFrame,最后一次性合并,能减少80%以上的内存碎片。
  • Parquet写入适配Hadoop:确保HDFS路径格式正确(比如hdfs://namenode:9000/user/data/output.parquet),如果你的环境已经配置好Hadoop客户端,也可以直接用本地路径映射的方式(比如/user/data/output.parquet)。
  • 超大数据量兜底方案:如果数据量突破千万级,内存中合并全量DataFrame还是有压力,可以改成先把每个批次的结果写入临时Parquet文件,最后再合并这些临时文件:
def parallel_process_with_temp_files(self, csv_path, temp_hdfs_dir, chunksize=10000):
    chunk_idx = 0
    for chunk in pd.read_csv(csv_path, chunksize=chunksize):
        processed_chunk = self.process_batch(chunk)
        # 写入HDFS临时目录
        temp_path = f"{temp_hdfs_dir}/temp_chunk_{chunk_idx}.parquet"
        processed_chunk.to_parquet(temp_path, engine='pyarrow', compression='snappy')
        chunk_idx += 1
    
    # 合并所有临时Parquet文件
    self.df = pd.read_parquet(f"{temp_hdfs_dir}/temp_chunk_*.parquet", engine='pyarrow')

这种方式完全不需要在内存中存储全量数据,适合处理超大文件。

内容的提问来源于stack exchange,提问作者N.J.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:15:26