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

如何基于Dask DataFrame的file列高效分文件写入CSV?

按指定列路径拆分写入Dask DataFrame到CSV文件

最优解决方案

通过先对file列进行shuffle分组,确保相同目标文件的行集中到同一个分区,再针对每个分区一次性写入对应文件,既降低内存消耗又提升效率:

import pandas as pd
import dask.dataframe as dd
from pathlib import Path

# 示例数据(替换为你的实际Dask DataFrame)
df = pd.DataFrame({"x": [1, 2, 3, 7, 11, 2], "y": [1, 1, 2, 8, 0, 0], "file": ["a.csv", "a.csv", "b.csv", "c.csv", "a.csv", "c.csv"]})
ddf = dd.from_pandas(df, npartitions=2)

# 关键步骤:按file列shuffle,让相同file的行归到同一分区
shuffled_ddf = ddf.shuffle(on="file", npartitions=len(ddf["file"].unique().compute()))

def write_single_file_partition(df: pd.DataFrame) -> None:
    # 每个分区仅对应一个目标文件
    target_file = df["file"].iloc[0]
    # 判断文件是否存在,决定是否写入表头
    write_header = not Path(target_file).exists()
    # 移除file列(可选,根据需求决定是否保留)
    output_df = df.drop(columns=["file"])
    output_df.to_csv(target_file, mode="a", header=write_header, index=False)

# 执行写入操作
shuffled_ddf.map_partitions(write_single_file_partition).compute()

原方案的问题分析

  1. 方案1(遍历唯一file):
    • 每次筛选单个file都会重新扫描整个Dask DataFrame,相当于重复读取数据N次(N为不同file的数量),IO成本极高,且多次compute会导致内存中堆积大量中间数据。
  2. 方案2(分区内分组append):
    • 同一目标文件会被多个分区重复打开/写入,引发IO竞争;同时Dask在compute时会保留所有分区的处理状态,无法及时释放内存,最终导致OOM。

额外优化建议

  • 调整分区数:如果不同file的数量较多,直接将npartitions设为唯一file的数量,让每个file对应一个独立分区,进一步降低单分区内存压力。
  • 禁用索引写入:始终设置index=False,避免写入不必要的索引列,减少数据量。
  • 配置Dask临时目录:通过dd.config.set({"temporary_directory": "/path/to/large/tmp"})让Dask将shuffle的临时数据写入磁盘,避免占用过多内存。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 04:45:22