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

如何解决合并大型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.4G5940345735
4.1G3140062568
25k5211
776M904023512
3.1G1674426150
1.7G1813465726
7.0G35947667115
848M685949713
681M685949711
6.3G34482809104
937M406520515
90M12604211
2.7M308301
7.1M1217651
611M569015310
627M569015310
5.9G2205256598
2.7G1558557744
3.4G1558557756
499M36919338
720M368367711
3.7G4216816560
3.3G3161965354

连接列基数

  • date: 1460+
  • att_1: 56
  • att_2: 3
  • att_3: 5
  • att_4: 5
  • att_5: 3
  • att_6: 3
  • att_7: 5

错误日志

执行dd.to_parquet时的错误序列:

  1. 大量shuffle警告:distributed.shuffle._scheduler_plugin - WARNING - Shuffle 99748e0284e4fcb627e68b744953791a
  2. 重复垃圾回收提示:distributed.utils_perf - WARNING - full garbage collections took 15% CPU time recently (threshold: 10%)
  3. 内存占用过高警告:WARNING - Unmanaged memory use is high. This may indicate a memory leak or the memory may not be released to the OS;
  4. Worker内存超限重启:WARNING - Worker tcp://127.0.0.1:36199 (pid=20601) exceeded 95% memory budget. Restarting...
  5. 最终错误: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 11:45:01