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

基于Dask的大CSV文件内存高效加载问题求助

可行解决方案

方案1:配置共享存储(最优推荐)

这是最符合Dask分布式设计的方案,不会占用客户端内存,加载过程完全分布式执行:

  • 将CSV文件上传到所有worker节点均可访问的共享存储介质,可选介质包括NFS共享目录、HDFS集群、私有对象存储(如MinIO)
  • 确保所有节点访问该文件的路径完全一致,比如NFS统一挂载到/data/share路径,文件路径为/data/share/cheque_data.csv
  • 直接使用dd.read_csv读取即可,代码示例:
import dask.dataframe as dd
from dask.distributed import Client

client = Client("tcp://你的scheduler地址:8786") # 替换为实际的scheduler地址
cheques = dd.read_csv("/data/share/cheque_data.csv")

方案2:客户端分块推送(无共享存储时可选)

如果暂时无法配置共享存储,可在客户端本地分块读取CSV后推送至集群,无需worker直接访问源文件:

import dask.dataframe as dd
import pandas as pd
from dask.distributed import Client
from dask import delayed

client = Client("tcp://你的scheduler地址:8786")
delayed_chunks = []
# 客户端侧分块读取,不会一次性加载全量数据到客户端内存
for chunk in pd.read_csv("cheque_data.csv", chunksize=10**4):
    delayed_chunks.append(delayed(chunk))
# 基于延迟分块构建Dask DataFrame
cheques = dd.from_delayed(delayed_chunks)
# 可选:将数据持久化到集群内存,避免后续操作重复读取传输
# cheques = cheques.persist()

原有Dask Bag方法的修复方案(不推荐)

你之前用Dask Bag报错的原因是pd.read_csv返回的分块是DataFrame对象,from_sequence会把单个DataFrame识别为一个元素而非一批行,调整后可以运行但不推荐:该方案需要将全量数据读取到客户端内存,大文件场景下会打爆客户端内存。

from dask.bag import from_sequence
import pandas as pd

all_records = []
for chunk in pd.read_csv("cheque_data.csv", chunksize=10**4):
    all_records.extend(chunk.to_dict("records"))
cheques = from_sequence(all_records).to_dataframe()

注意事项

  • 方案2的传输效率受客户端与集群的网络带宽限制,超过10GB的文件优先选择方案1
  • 若使用对象存储,Dask原生支持S3协议,直接传入s3://bucket/cheque_data.csv格式的路径即可读取,无需额外挂载

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 15:27:03