如何在Pandas中并行读取多个CSV文件至DataFrame?
Pandas并行读取带注释CSV文件报错问题解决
问题背景
有一批CSV文件,每个文件开头包含空行和注释行,后续是foo、bar、baz三列的数值数据,同时存在ID与文件路径的映射字典(如files = {'A': 'A.csv', 'B': 'B.csv'})。
串行读取的代码可正常运行:
import pandas as pd columns = ['foo', 'bar', 'baz'] skip = 4 df = (pd.concat({k: pd.read_csv(v, skiprows=skip, sep=r'\s+', names=columns) for k,v in files.items()}, names=['ID']) .reset_index('ID') .reset_index(drop=True) )
尝试用joblib实现并行读取时,出现报错:
from joblib import Parallel, delayed from multiprocessing import cpu_count from pathlib import Path n_jobs = cpu_count() def read_file(res_dict: dict, skiprows: int, columns: list[str], id: str, file: Path ) -> None: res_dict[id] = pd.read_csv(file, skiprows=skiprows, sep=r'\s+', names=columns) temp = {} temp = Parallel(n_jobs)(delayed(read_file)(temp, skip_rows, columns, id, file) for id, file in master2file.items()) df = (pd.concat(temp, names=['ID']) .reset_index('ID') .reset_index(drop=True) )
报错信息:
Traceback (most recent call last):
File "/home/...py", line 54, in
df = (pd.concat(temp,
File "/home/../.venv/lib/python3.10/site-packages/pandas/core/reshape/concat.py", line 372, in concat
op = _Concatenator(
File "/home/../.venv/lib/python3.10/site-packages/pandas/core/reshape/concat.py", line 452, in init
raise ValueError("All objects passed were None")
ValueError: All objects passed were None
错误原因
- 函数返回值为空:
read_file函数返回None,joblib的Parallel会收集每个任务的返回值,最终temp列表全是None,导致pd.concat无法找到有效数据。 - 多进程内存隔离:传递给子进程的
temp字典是父进程的副本,子进程修改的是自己内存空间里的副本,父进程的temp不会被更新,实际没有收集到任何DataFrame。 - 参数名不一致:代码中使用了未定义的
skip_rows,而之前定义的变量是skip,属于笔误。
修正后的并行实现
核心思路是让并行任务直接返回(ID, DataFrame)的元组,通过joblib收集返回值后再构建字典进行拼接:
from joblib import Parallel, delayed from multiprocessing import cpu_count import pandas as pd from pathlib import Path # 定义参数 columns = ['foo', 'bar', 'baz'] skip = 4 master2file = {'A': Path('A.csv'), 'B': Path('B.csv')} # ID与文件的映射字典 def read_file(id: str, file: Path, skiprows: int, columns: list[str]) -> tuple[str, pd.DataFrame]: # 读取文件并返回ID和对应的DataFrame df = pd.read_csv(file, skiprows=skiprows, sep=r'\s+', names=columns) return (id, df) # 并行执行任务 n_jobs = cpu_count() results = Parallel(n_jobs=n_jobs)( delayed(read_file)(id, file, skip, columns) for id, file in master2file.items() ) # 将结果转换为字典后拼接 df = (pd.concat({k: v for k, v in results}, names=['ID']) .reset_index('ID') .reset_index(drop=True) )
关键修正点
- 避免共享可变对象:多进程环境下,不同进程拥有独立内存空间,通过返回值传递数据是可靠的方式,而非修改共享字典。
- 函数返回有效数据:让
read_file返回(ID, DataFrame)元组,确保joblib能收集到可用的数据。 - 修正参数名错误:使用定义好的
skip变量替代未定义的skip_rows。
额外优化建议
- 可以添加
verbose参数(如verbose=10)到Parallel中,查看并行任务的执行进度。 - 若运行在Windows系统,需将主逻辑放在
if __name__ == '__main__':代码块中,避免多进程重复初始化问题:if __name__ == '__main__': results = Parallel(n_jobs=n_jobs)(...) df = pd.concat(...) - 对于超大文件,可考虑在读取时指定
dtype参数减少内存占用,提升读取速度。
内容的提问来源于stack exchange,提问作者DeltaIV
相关产品推荐
相关产品推荐

