使用Dask DataFrame处理大CSV文件时遇AssertionError报错,求解决及性能优化建议
Dask DataFrame处理大CSV文件时遇AssertionError报错,求解决及性能优化建议
嘿,我来帮你拆解下当前的问题,以及优化你的Dask处理流程~
一、AssertionError的报错原因
你代码里的核心问题出在map_partitions的调用逻辑上:
chunk1_df = chunk_df.map_partitions(chunk_df,barra).compute()
map_partitions的第一个参数要求是可调用的函数(也就是你写的Proceso函数),但你却把chunk_df(一个Dask DataFrame对象)传了进去——这直接触发了assert callable(func)的断言失败,因为Dask DataFrame并不是可调用的函数类型。
二、修正后的完整Dask代码(更高效简洁)
而且你完全不需要用pandas.read_csv手动分块再转Dask,Dask本身就专为大文件并行处理设计,能自动分块、利用多核CPU,代码会简洁很多,效率也更高:
import dask.dataframe as dd from datetime import datetime from tkinter import filedialog as fd dt1 = datetime.now() # 定义需要筛选的barra列表 barra_list = ['barra1','barra3','barra5','barra7','barra9','barra11','barra13','barra15','barra17','barra19'] # 选择输入输出路径 filename1 = fd.askopenfilename(title="ARCHIVO CSV ENTRADA") output_path = fd.asksaveasfilename(initialfile='Untitled.csv', defaultextension=".csv", title="ARCHIVO DE SALIDA") # 1. 直接用Dask读取大CSV,提前指定列名(避免在每个分区重复重命名) # 注意:如果你的CSV第一行是表头,把header=0加上;如果无表头用header=None # 请根据原文件实际列数调整names!原描述是10列,这里对应设置列名 df = dd.read_csv( filename1, names=['barra', 'D1', 'D2','D3','D4','D5','D6','D7','D8','CMG'], header=None, # 原文件有表头则改为header=0 blocksize="64MB", # 可根据内存调整,默认64MB,内存充足可设为128/256MB dtype={'barra': 'string'} # 显式指定类型,减少内存占用和推断开销 ) # 2. 直接筛选符合条件的行,Dask自动并行处理所有分区 filtered_df = df[df['barra'].isin(barra_list)] # 3. 写入结果CSV,Dask自动处理分块写入,无需手动append filtered_df.to_csv(output_path, index=False, single_file=True) # 计算总耗时 dt2 = datetime.now() print(f"FIN: {dt2-dt1}")
如果你需要保留自定义处理函数(比如后续扩展复杂逻辑)
如果你的Proceso函数后续要添加更复杂的处理逻辑,正确的map_partitions调用方式是这样的:
# 定义自定义处理函数 def Proceso(chunk_df, target_barras): # 这里可以添加更复杂的分区内处理逻辑 chunk1 = chunk_df[chunk_df['barra'].isin(target_barras)] return chunk1 # 读取文件(和上面一致) df = dd.read_csv( filename1, names=['barra', 'D1', 'D2','D3','D4','D5','D6','D7','D8','CMG'], header=None, dtype={'barra': 'string'} ) # 正确调用map_partitions:第一个参数是函数,后面是传给函数的额外参数 filtered_df = df.map_partitions(Proceso, barra_list, meta=df.dtypes)
这里的meta参数是告诉Dask输出数据的列类型,避免自动推断的额外开销。
三、你的其他疑问解答
能不能去掉pandas的chunk部分?
完全可以,而且非常推荐!Dask的dd.read_csv会自动将大文件拆分为多个分区并行处理,比手动用pandas.read_csv(chunksize=...)循环处理高效得多——后者是单线程的,而Dask能利用你的多核CPU并行计算。性能优化额外建议
- 批量处理多文件:你有80个32GB的CSV,可以直接用
dd.read_csv("path/to/your/*.csv")读取所有文件,Dask会自动合并处理,不用逐个文件手动操作。 - 切换并行后端:如果你的处理是CPU密集型的,可切换到多进程后端提升效率:
from dask.distributed import Client client = Client() # 启动本地多进程集群 - 减少内存占用:尽可能给所有列指定
dtype(比如数值列设为float32而非float64,字符串列设为string而非object),大幅降低内存压力。
备注:内容来源于stack exchange,提问作者Eduardo
相关产品推荐
相关产品推荐

