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

如何并行化Pandas的.groupby(...).size()操作?适配大DF与难序列化类

如何并行化Pandas的groupby.size()操作?

我在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 16:05:00