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

如何在内存受限的Windows 10笔记本上完成Dask DataFrame的groupby-size任务?

解决Dask DataFrame Groupby-Size内存不足的方案

我之前在类似的资源受限环境下处理过超大数据集的分组计数任务,结合你的场景(12GB内存Windows笔记本、7GB→21GB Parquet、3亿条记录),给你几个可行的方案,按推荐优先级排序:

1. 先局部计数再全局聚合(最优方案)

groupby.size()默认会触发全量shuffle,这是内存压力的主要来源。我们可以先在每个分区上做局部分组计数,再把所有局部结果合并求和,这样shuffle的数据量会大幅减少:

import dask.dataframe as dd

def partial_group_count(df):
    # 对单个分区的col_1、col_2分组计数
    return df.groupby(["col_1", "col_2"]).size()

# 读取Parquet文件(保持原有分区即可)
ddf = dd.read_parquet(parquet_path)

# 生成每个分区的局部计数结果
partial_counts = ddf.map_partitions(
    partial_group_count,
    meta=("size", "int64")  # 指定输出元数据类型,避免Dask推断出错
)

# 对所有局部结果按分组键求和,得到最终计数
final_counts = partial_counts.groupby(partial_counts.index).sum()

# 导出结果
final_counts.to_csv(csv_path)

这个方法的核心是把大数据集的全局shuffle转化为小数据集的聚合,内存占用能降到原来的几十分之一,同时速度也不会太慢。

2. 调整Dask分区与内存配置

你的Parquet分片是235MB/个,对于12GB内存来说还是偏大,而且Dask默认的内存管理可能没有适配你的环境:

步骤1:缩小分区大小

把每个大分片拆分成更小的分区,让单个分区的处理内存控制在系统能承受的范围内:

ddf = dd.read_parquet(parquet_path).repartition(partition_size="64MB")

64MB是一个比较保守的数值,你也可以根据实际内存使用调整到96MB或128MB。

步骤2:配置Dask内存限制与溢出策略

通过Dask配置强制限制内存使用,并允许溢出到磁盘:

import dask

dask.config.set({
    "memory.limit": "8GB",  # 给Dask分配8GB内存,留4GB给Windows系统和其他进程
    "memory.spill": True,  # 内存不足时自动spill到磁盘
    "worker.memory.target": 0.6,  # 内存使用到60%就开始准备spill
    "worker.memory.terminate": 0.95  # 内存用到95%才终止进程,减少崩溃概率
})

# 之后执行原代码
ddf = dd.read_parquet(parquet_path).repartition(partition_size="64MB")
sr = ddf.groupby(["col_1", "col_2"]).size()
sr.to_csv(csv_path)

步骤3:用本地集群精细控制资源

如果是多进程模式,多个Worker会抢占内存,建议手动创建本地集群,限制Worker数量和单个Worker的内存:

from dask.distributed import Client

# 创建2个Worker,每个最多用4GB内存,总占用8GB
client = Client(n_workers=2, memory_limit="4GB")

ddf = dd.read_parquet(parquet_path).repartition(partition_size="64MB")
sr = ddf.groupby(["col_1", "col_2"]).size()
sr.to_csv(csv_path)

client.close()  # 任务完成后关闭集群

3. 分批次处理数据(极端内存紧张时用)

如果上面的方法还是不行,可以按col_1的唯一值拆分数据,逐批次处理后合并结果:

import dask.dataframe as dd
import pandas as pd

ddf = dd.read_parquet(parquet_path)

# 先获取所有唯一的col_1值(这一步内存占用很小)
unique_col1_values = ddf["col_1"].unique().compute()

all_results = []
for val in unique_col1_values:
    # 过滤出当前col_1对应的数据子集
    subset = ddf[ddf["col_1"] == val]
    # 对col_2分组计数并计算结果
    batch_counts = subset.groupby("col_2").size().compute()
    # 把col_1值补回到结果中
    batch_df = batch_counts.reset_index()
    batch_df["col_1"] = val
    all_results.append(batch_df)

# 合并所有批次的结果,做最终求和
final_df = pd.concat(all_results).groupby(["col_1", "col_2"]).sum()
final_df.to_csv(csv_path)

这个方法内存占用最低,但如果col_1的唯一值很多,速度会比较慢,适合极端内存不足的场景。

额外提示

  • 处理Parquet时,确保你用的是最新版的Dask和PyArrow/fastparquet,新版本对内存优化更好;
  • Windows系统下关闭不必要的后台程序,释放更多可用内存;
  • 如果后续数据扩大到21GB,方案1和方案2的适配性最好,方案3可能需要更长时间。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 03:44:07