如何在内存受限的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
相关产品推荐
相关产品推荐

