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

使用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输出数据的列类型,避免自动推断的额外开销。

三、你的其他疑问解答

  1. 能不能去掉pandas的chunk部分?
    完全可以,而且非常推荐!Dask的dd.read_csv会自动将大文件拆分为多个分区并行处理,比手动用pandas.read_csv(chunksize=...)循环处理高效得多——后者是单线程的,而Dask能利用你的多核CPU并行计算。

  2. 性能优化额外建议

  • 批量处理多文件:你有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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 17:58:11