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

如何分分区迭代创建Dask DataFrame并在其未完成时逐步写入磁盘?

迭代创建Dask DataFrame并高效拼接的实现方案

针对内存不足的问题,核心思路是利用Dask的**延迟执行(delayed)**特性,避免将所有处理后的pandas DataFrame同时驻留内存。具体实现步骤如下:

改造后的代码示例

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

def process_data(data) -> pd.DataFrame:
    # 原有的数据处理逻辑保持不变
    return df

# 生成延迟执行的任务列表,不立即计算
delayed_dfs = []
for data in datasets:
    # 用delayed包装处理函数,返回延迟对象而非实际的pandas DF
    delayed_df = delayed(process_data)(data)
    delayed_dfs.append(delayed_df)

# 将延迟对象列表转换为Dask DataFrame
final_dask_df = dd.from_delayed(delayed_dfs)

# 保存结果到CSV,Dask会自动分块处理,无需一次性加载所有数据
final_dask_df.to_csv("result_*.csv", index=False)
# 如果需要单个CSV文件(注意:大文件可能仍有内存压力),可以添加参数single_file=True
# final_dask_df.to_csv("result.csv", index=False, single_file=True)

关键说明

  • 延迟执行:delayed(process_data)(data)不会立即执行处理逻辑,而是生成一个任务对象记录要执行的操作。直到调用to_csv或compute()时,Dask才会调度任务,分批次处理数据集。
  • 结构一致性:确保process_data返回的所有pandas DataFrame具有相同的列名和数据类型,否则dd.from_delayed会报错。如果存在结构差异,可以在处理函数中统一列结构。
  • 内存优化:Dask会根据系统内存自动调整并行处理的任务数,每次仅将部分数据加载到内存处理,完成后立即写入磁盘,避免内存溢出。
  • 并行提速:默认情况下Dask会利用多核CPU并行处理数据集,比单线程的pandas循环更快。如果需要控制并行度,可以通过dask.config.set(scheduler='processes', num_workers=4)调整。

与原pandas代码的差异

原代码会将所有处理后的pandas DataFrame存入列表,最终一次性拼接,内存占用等于所有结果的总和;而改造后的Dask方案仅在需要时处理单个数据集,内存占用仅为单个(或少量)数据集处理后的大小。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 23:35:20