Spark按id分组统计唯一grade、求最大值等聚合操作实现
Spark 聚合计算实现
现有测试数据集,需基于Spark完成三项聚合计算逻辑,以下是完整实现步骤。
测试数据准备
首先执行以下代码创建测试DataFrame:
from pyspark.sql import functions as F # 定义表结构与测试数据 columns = ['id', 'grade', 'date'] values = [('101','good','2022/06/01'), ('102','good','2022/06/01'), ('103','fail','2022/06/02'), ('104','poor','2022/06/02'),('101','good','2022/06/08'), ('101','excellent','2022/06/14'),('102','poor','2022/06/10'), ('104','good','2022/06/09'),('102','poor','2022/06/13'), ('103','fail','2022/06/14')] # 创建DataFrame DataFrame_from_Scratch = spark.createDataFrame(values, columns) df = DataFrame_from_Scratch df.show()
需求实现
1. 按id分组统计每个id对应的唯一grade数量
使用groupBy按id分组,搭配countDistinct函数对grade列去重计数即可,代码如下:
id_grade_distinct_cnt = df.groupBy("id") \ .agg(F.countDistinct("grade").alias("distinct_grade_count")) id_grade_distinct_cnt.show()
执行输出结果:
+---+--------------------+ | id|distinct_grade_count| +---+--------------------+ |101| 2| |102| 2| |103| 1| |104| 2| +---+--------------------+
2. 获取上述去重计数结果的最大值
基于第一步得到的统计结果,对计数字段取最大值即可:
max_cnt = id_grade_distinct_cnt \ .agg(F.max("distinct_grade_count").alias("max_count")) \ .collect()[0]["max_count"] print(max_cnt)
执行输出:2
3. 按id、date字段分组,统计对应的唯一date数量
按id和date双字段分组后,每组内date值完全一致,直接对date列去重计数即可:
id_date_distinct_cnt = df.groupBy("id", "date") \ .agg(F.countDistinct("date").alias("distinct_date_count")) id_date_distinct_cnt.show()
执行输出结果:
+---+----------+-------------------+ | id| date|distinct_date_count| +---+----------+-------------------+ |101|2022/06/01| 1| |101|2022/06/08| 1| |101|2022/06/14| 1| |102|2022/06/01| 1| |102|2022/06/10| 1| |102|2022/06/13| 1| |103|2022/06/02| 1| |103|2022/06/14| 1| |104|2022/06/02| 1| |104|2022/06/09| 1| +---+----------+-------------------+
内容的提问来源于stack exchange,提问作者umagba alex
相关产品推荐
相关产品推荐

