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

并行读取CSV并合并DataFrame的代码问题及优化咨询

嘿,咱们一步步来解决你的并行CSV处理问题吧!我会先拆解你当前代码里的问题,再给出优化方案,让这个批量处理任务更顺畅高效。

批量CSV并行读取与合并的问题解决

原代码的核心问题拆解

  1. 回调函数无法区分任务来源:你的concatDf没有接收回调返回的结果,还硬编码要合并df_1和df_2,完全没法适配动态的多文件任务结果。
  2. 无法判断所有任务完成:原代码没有等待所有线程任务结束,也没有统一收集所有返回的DataFrame,根本没法合并50个文件的结果。
  3. 实现方式效率低下: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会自动管理并行任务,不需要手动维护线程池,非常适合超大数据场景。


总结

  1. 放弃回调函数,用ThreadPoolExecutor.map批量收集任务结果,天然解决任务区分和完成判断的问题。
  2. 优化pd.read_csv的参数,减少内存占用和读取时间。
  3. 内存压力大时用分阶段合并,超大数据场景直接用Dask。

内容的提问来源于stack exchange,提问作者Invictus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 00:27:47