如何通过Databricks API(SQL/PySpark)实现Pandas的groupby聚合映射原表功能
在Databricks中实现Pandas
transform 分组聚合并添加新列的功能 对应Pandas中通过groupby.transform计算分组统计量并追加回原DataFrame的需求,在Databricks里可以用PySpark或者SQL两种方式实现,具体如下:
方法一:PySpark 实现
利用Spark的窗口函数(Window)完成分组统计的广播映射,和Pandas的transform逻辑完全一致:
步骤1:创建示例DataFrame
from pyspark.sql.window import Window import pyspark.sql.functions as F # Databricks环境中SparkSession通常已自动初始化,无需额外创建 data = [("f1", "jack", 12), ("f2", "jen", 23), ("f3", "joe", 13), ("f1", "jan", 15)] df = spark.createDataFrame(data, schema=["class", "user", "screen"])
步骤2:定义窗口规范并计算分组标准差
# 按class列分组的窗口规则 window_spec = Window.partitionBy("class") # 添加分组标准差列gp,匹配Pandas的样本标准差逻辑 df_with_gp = df.withColumn("gp", F.stddev("screen").over(window_spec)) # 查看结果 df_with_gp.show()
输出结果
+-----+----+------+------------------+ |class|user|screen| gp| +-----+----+------+------------------+ | f1|jack| 12|2.1213203435596424| | f1| jan| 15|2.1213203435596424| | f2| jen| 23| null| | f3| joe| 13| null| +-----+----+------+------------------+
方法二:Databricks SQL 实现
先将数据注册为临时视图,再通过SQL窗口函数完成计算:
步骤1:注册临时视图
# 在PySpark中创建临时视图,供SQL查询调用 df.createOrReplaceTempView("user_screen_data")
步骤2:执行SQL查询
在Databricks SQL编辑器或PySpark中执行以下语句:
SELECT class, user, screen, STDDEV(screen) OVER (PARTITION BY class) AS gp FROM user_screen_data
输出结果
和PySpark实现的结果完全一致,单行分组的gp列会显示为NULL(对应Pandas的NaN)。
补充说明
- Spark的
STDDEV()默认计算样本标准差,和Pandasstd()默认的ddof=1逻辑一致;若需总体标准差,可改用STDDEV_POP()函数。 - 窗口函数
OVER (PARTITION BY class)会自动将分组计算结果映射回原数据的每一行,完全复现Pandastransform的效果。
内容的提问来源于stack exchange,提问作者Leon Adams
相关产品推荐
相关产品推荐

