如何解决合并大型Dask DataFrame时的内存错误?
问题:Dask合并多CSV转Parquet时内存不足导致Worker崩溃
尝试将23个CSV读取为Dask DataFrame,合并后导出Parquet时因内存问题失败。此前用Pandas处理,数据量增大后仅合并少数文件就触发内核OOM(通过dmesg确认:Out of memory: Killed process 25693 (ld-linux-x86-64) total-vm:125749916kB)。
已尝试但无效的方法
- 基于
date列建立索引并排序连接 - 以
date为索引,用df.repartition(freq='7d')调整分区,修改频率后仍出现worker died unexpectedly错误 - 仅将基础表设为Dask DataFrame,其余用Pandas读取,因任务图达49GiB触发序列化错误
- 开启shuffle分区压缩(snappy),无效果
- 注:逐个将22个CSV与基础DataFrame连接是可行的,说明批量合并是问题根源
代码结构
from dask.distributed import LocalCluster import dask.dataframe as dd import dask cluster = LocalCluster() dask_client = cluster.get_client() dask.config.set({'distributed.worker.shuffle-compression': 'snappy'}) base_df = dd.read_csv( my_path, keep_default_na=False, parse_dates=[6], date_format="%Y-%m-%d", blocksize='100MB' ) base_df = base_df.set_index("date") file_list = ['file_1','file_2','file_3',...,'file_22'] for file_name in file_list: metrics_file = "/my/path/" + file_name + '.csv' metric_df = dd.read_csv( metrics_file, keep_default_na=False, parse_dates=[6], date_format="%Y-%m-%d" ) metric_df = metric_df.set_index("date") base_df = dd.merge(base_df, metric_df, on=['date','att_1','att_2','att_3','att_4','att_5','att_6','att_7'], how="left") dd.to_parquet( base_df, 's3://my-bucket/path', engine='pyarrow', compression='snappy', overwrite=True, write_index=False, storage_options={'key': aws.access_key, 'secret': aws.secret_key} )
文件特征
| 大小 | 行数 | 分区数 |
|---|---|---|
| 3.4G | 59403457 | 35 |
| 4.1G | 31400625 | 68 |
| 25k | 521 | 1 |
| 776M | 9040235 | 12 |
| 3.1G | 16744261 | 50 |
| 1.7G | 18134657 | 26 |
| 7.0G | 35947667 | 115 |
| 848M | 6859497 | 13 |
| 681M | 6859497 | 11 |
| 6.3G | 34482809 | 104 |
| 937M | 4065205 | 15 |
| 90M | 1260421 | 1 |
| 2.7M | 30830 | 1 |
| 7.1M | 121765 | 1 |
| 611M | 5690153 | 10 |
| 627M | 5690153 | 10 |
| 5.9G | 22052565 | 98 |
| 2.7G | 15585577 | 44 |
| 3.4G | 15585577 | 56 |
| 499M | 3691933 | 8 |
| 720M | 3683677 | 11 |
| 3.7G | 42168165 | 60 |
| 3.3G | 31619653 | 54 |
连接列基数
date: 1460+att_1: 56att_2: 3att_3: 5att_4: 5att_5: 3att_6: 3att_7: 5
错误日志
执行dd.to_parquet时的错误序列:
- 大量shuffle警告:
distributed.shuffle._scheduler_plugin - WARNING - Shuffle 99748e0284e4fcb627e68b744953791a - 重复垃圾回收提示:
distributed.utils_perf - WARNING - full garbage collections took 15% CPU time recently (threshold: 10%) - 内存占用过高警告:
WARNING - Unmanaged memory use is high. This may indicate a memory leak or the memory may not be released to the OS; - Worker内存超限重启:
WARNING - Worker tcp://127.0.0.1:36199 (pid=20601) exceeded 95% memory budget. Restarting... - 最终错误:
Attempted to run task on 4 different workers, but all those workers died while running it.
主机配置
- OS: Linux
- 内存: 124G
- 磁盘: 493G
- CPU: 16核(32线程)
- 架构: x86_64
解决方案建议
1. 优化合并策略:切断任务图链式依赖
当前循环中每次合并都会累积任务图,22次合并后任务图会指数级膨胀,导致调度和内存压力剧增。改为每次合并后立即持久化(persist),将中间结果落地,释放历史任务依赖:
base_df = base_df.set_index("date").persist() # 先持久化基础表 for file_name in file_list: metrics_file = "/my/path/" + file_name + '.csv' metric_df = dd.read_csv( metrics_file, keep_default_na=False, parse_dates=[6], date_format="%Y-%m-%d" ) metric_df = metric_df.set_index("date").persist() # 持久化当前要合并的表 base_df = dd.merge(base_df, metric_df, on=['date','att_1','att_2','att_3','att_4','att_5','att_6','att_7'], how="left") base_df = base_df.persist() # 合并后立即持久化,压缩任务图
2. 调整分区策略:基于完整连接键分区
当前仅按date分区,但合并键包含多个字段,Dask需要将相同连接键的数据放到同一分区才能避免大量跨分区传输。可以:
- 用连接键的哈希值作为分区键,确保同键数据在同一分区:
def get_partition_key(df): return hash(tuple(df[['date','att_1','att_2','att_3','att_4','att_5','att_6','att_7']].values)) # 对基础表和所有metric表执行分区 base_df = base_df.set_index(base_df.apply(get_partition_key, meta=('int', int)).compute())
- 或者用
partition_size控制分区大小(建议100-200MB):
base_df = base_df.repartition(partition_size='150MB')
3. 限制Worker内存并开启磁盘溢出
手动配置LocalCluster的Worker资源,避免内存耗尽:
cluster = LocalCluster( memory_limit='20GB', # 每个Worker分配20GB,124G可设5个Worker n_workers=5, threads_per_worker=6 ) # 开启内存溢出到磁盘,减少OOM概率 dask.config.set({ 'distributed.worker.memory.spill': True, 'distributed.worker.memory.target': 0.6, # 内存使用到60%开始溢出 'distributed.worker.memory.terminate': 0.95 })
4. 优化CSV读取
- 明确指定
dtype,避免Dask自动推断类型消耗内存:
dtype_dict = { 'att_1': 'int32', 'att_2': 'int8', 'att_3': 'int8', # 其他列按需定义类型 } base_df = dd.read_csv( my_path, keep_default_na=False, parse_dates=[6], date_format="%Y-%m-%d", blocksize='100MB', dtype=dtype_dict )
- 小文件直接用Pandas读取后转Dask,减少分区开销:
import pandas as pd import os for file_name in file_list: metrics_file = "/my/path/" + file_name + '.csv' if os.path.getsize(metrics_file) < 100*1024*1024: # 小于100MB的文件 pd_df = pd.read_csv(metrics_file, keep_default_na=False, parse_dates=[6], date_format="%Y-%m-%d") metric_df = dd.from_pandas(pd_df, npartitions=1) else: metric_df = dd.read_csv(metrics_file, keep_default_na=False, parse_dates=[6], date_format="%Y-%m-%d")
5. 导出前重新分区
合并后的DataFrame可能分区数量过多或分布不均,导出前重新调整分区:
base_df = base_df.repartition(partition_size='200MB') dd.to_parquet(base_df, ...)
内容的提问来源于stack exchange,提问作者ifightfortheuserz
相关产品推荐
相关产品推荐

