如何在Spark DataFrame中对分组记录执行自定义计算?
按"foo"和"bar"分组并自定义生成"badger"字段的实现方案
先把你的原始DataFrame用pandas构造出来,方便后续演示:
import pandas as pd df = pd.DataFrame({ 'foo': [1, 1, 2, 3, 3], 'bar': [2, 2, 2, 4, 4], 'fox': ['red', 'red', 'brown', 'taupe', 'red'], 'cow': ['blue', 'yellow', 'green', 'fuschia', 'orange'] })
接下来,核心思路是用groupby(['foo', 'bar'])分组,然后通过apply()方法传入自定义函数——因为apply()允许你对每个分组做任意逻辑处理,包括增删记录、生成新字段,完全满足你的需求。
示例自定义逻辑
假设我们的自定义规则是:
- 对于每个分组:
- 如果分组内
fox有多种不同取值,生成两条记录:- 第一条:
badger为fox所有唯一值拼接 + '-' + cow所有唯一值拼接 - 第二条:
badger为fox的众数 + '-' + cow的众数
- 第一条:
- 如果分组内
fox取值唯一,仅生成一条记录:badger为fox值 + '-' + cow所有值拼接
- 如果分组内
对应的自定义函数和调用代码如下:
def custom_calculate(group): # 获取当前分组的foo和bar值(每个分组内这两个值是固定的) foo_val = group['foo'].iloc[0] bar_val = group['bar'].iloc[0] # 提取fox和cow的相关信息 fox_unique = ', '.join(group['fox'].unique()) cow_unique = ', '.join(group['cow'].unique()) fox_mode = group['fox'].mode()[0] cow_mode = group['cow'].mode()[0] cow_all = ', '.join(group['cow']) # 根据规则生成结果 result_rows = [] if len(group['fox'].unique()) > 1: # 生成两条记录 result_rows.append({ 'foo': foo_val, 'bar': bar_val, 'badger': f"{fox_unique} - {cow_unique}" }) result_rows.append({ 'foo': foo_val, 'bar': bar_val, 'badger': f"{fox_mode} - {cow_mode}" }) else: # 生成一条记录 result_rows.append({ 'foo': foo_val, 'bar': bar_val, 'badger': f"{group['fox'].iloc[0]} - {cow_all}" }) # 返回当前分组处理后的DataFrame return pd.DataFrame(result_rows) # 分组并应用自定义函数,最后合并所有分组的结果 final_df = df.groupby(['foo', 'bar']).apply(custom_calculate).reset_index(drop=True) print(final_df)
运行结果
执行后会得到这样的输出:
foo bar badger 0 1 2 red - blue, yellow 1 2 2 brown - green 2 3 4 taupe, red - fuschia, orange 3 3 4 red - orange
可以看到,分组1-2(fox值唯一)只保留了1条记录;分组3-4(fox有两个不同值)生成了2条记录,完全实现了自定义增删记录和计算的需求。
你只需要根据自己的实际业务逻辑修改custom_calculate函数里的规则即可——不管是复杂的统计、过滤还是生成新记录,这个模式都能适配。
内容的提问来源于stack exchange,提问作者Sasgorilla
相关产品推荐
相关产品推荐

