如何并行化Pandas的.groupby(...).size()操作?适配大DF与难序列化类
我在Python 3.10的类方法(非__init__)里写了这段代码:
self.features = self.features.groupby(["token", "feature"], as_index=False).size() \ .rename(columns={"size": "freq"})
self.features是处理大量文本生成的大型DataFrame,其中包含难以序列化的自定义类元素——我已经在用dill做序列化,并行任务也换成了pathos替代标准多进程。
请问有没有办法并行化上面这段groupby(...).size()的处理?我知道一些Pandas并行化方法,但它们大多依赖速度较慢的.apply(),不想用这种方案。
可行的并行化方案
1. 基于pathos的分块聚合
既然你已经在用pathos,可以把大型DataFrame拆分成多个小分块,每个分块单独执行groupby.size(),最后再把所有分块的结果合并,做一次全局聚合得到最终频次。这种方式避免了序列化整个大DataFrame,pathos的多进程也能兼容dill处理自定义类。
示例代码:
from pathos.pools import ProcessPool def process_chunk(chunk): return chunk.groupby(["token", "feature"], as_index=False).size() # 把DataFrame拆分成N个块,N建议和CPU核心数匹配 chunks = [self.features.iloc[i:i+100000] for i in range(0, len(self.features), 100000)] # 用pathos的进程池并行处理 with ProcessPool() as pool: chunk_results = pool.map(process_chunk, chunks) # 合并所有分块结果,再做一次groupby求和得到最终频次 final_result = pd.concat(chunk_results) \ .groupby(["token", "feature"], as_index=False)["size"].sum() \ .rename(columns={"size": "freq"}) self.features = final_result
注意:分块大小可以根据你的内存情况调整,避免单个分块占用过多内存。
2. 用Dask实现自动并行
Dask是专门做并行数据分析的库,它可以直接处理Pandas的DataFrame,自动将任务拆分成并行执行的子任务,而且可以配置使用dill来序列化自定义对象。
示例代码:
import dask.dataframe as dd from dask.config import set_options # 配置Dask使用dill序列化,兼容自定义类 set_options({"serialization": "dill"}) # 将Pandas DataFrame转为Dask DataFrame dask_df = dd.from_pandas(self.features, npartitions=8) # npartitions建议设为CPU核心数的2-4倍 # 执行groupby.size(),然后转回Pandas DataFrame final_result = dask_df.groupby(["token", "feature"]).size() \ .reset_index(name="freq") \ .compute() self.features = final_result
这种方式不需要手动拆分数据,Dask会自动处理并行逻辑,适合不想写太多拆分代码的场景。
3. 提取分组键后并行统计频次
groupby.size()本质是统计分组键的出现次数,你可以单独提取token和feature列,转成可哈希的元组列表,然后用pathos的进程池并行统计频次,最后再转回DataFrame。这种方式只需要序列化分组键,不需要处理整个DataFrame,效率更高。
示例代码:
from pathos.pools import ProcessPool from collections import defaultdict def count_partition(partition): counter = defaultdict(int) for token, feature in partition: counter[(token, feature)] += 1 return counter # 提取分组键为元组列表 group_keys = list(zip(self.features["token"], self.features["feature"])) # 拆分分组键为多个分区 partition_size = len(group_keys) // 8 # 8为进程数,可根据CPU核心调整 key_partitions = [group_keys[i:i+partition_size] for i in range(0, len(group_keys), partition_size)] # 并行统计 with ProcessPool() as pool: partition_counters = pool.map(count_partition, key_partitions) # 合并所有计数器 final_counter = defaultdict(int) for counter in partition_counters: for k, v in counter.items(): final_counter[k] += v # 转成DataFrame final_result = pd.DataFrame(final_counter.items(), columns=["token_feature", "freq"]) final_result[["token", "feature"]] = pd.DataFrame(final_result["token_feature"].tolist(), index=final_result.index) final_result = final_result.drop("token_feature", axis=1)[["token", "feature", "freq"]] self.features = final_result
这种方式的优势是内存占用更低,只处理分组键,适合DataFrame中其他列数据量极大的场景。
内容的提问来源于stack exchange,提问作者TiMauzi

