Spark DataFrame添加分组最大值标记列的实现问题求助
Spark DataFrame添加分组最大值标记列的实现问题求助
我最近在处理一个包含70多列的Spark DataFrame,里面有number和pk两列,想要新增一列highest_check,用来标记每一行的number是否是其对应pk分组里的最大值——如果是就设为True,否则设为False。
为了简化问题,我先拿一个只包含number和pk的示例DataFrame dm1来举例,数据如下:
| number | pk |
|---|---|
| 100 | k1 |
| 200 | k2 |
| 300 | k3 |
| 400 | k4 |
| 330 | k3 |
| 500 | k5 |
| 220 | k2 |
| 370 | k5 |
| 180 | k2 |
| 90 | k1 |
| 470 | k4 |
我期望得到的最终结果是这样的:
| number | pk | highest_check |
|---|---|---|
| 100 | k1 | True |
| 200 | k2 | False |
| 300 | k3 | False |
| 400 | k4 | False |
| 330 | k3 | True |
| 500 | k5 | True |
| 220 | k2 | True |
| 370 | k5 | False |
| 180 | k2 | False |
| 90 | k1 | False |
| 470 | k4 | True |
我自己的尝试及遇到的问题
一开始我想先生成一个包含各pk分组最大值的DataFrame dm2,再通过匹配来给原表标记。我写的代码是这样的:
from pyspark.sql import functions as F from pyspark.sql.window import Window w = Window.partitionBy('pk') dm2 = dm1.withColumn('maxnumber', F.max('number').over(w))\ .where(F.col('number') == F.col('maxnumber'))
生成的dm2数据如下:
| number | pk | maxnumber |
|---|---|---|
| 100 | k1 | 100 |
| 220 | k2 | 220 |
| 330 | k3 | 330 |
| 470 | k4 | 470 |
| 500 | k5 | 500 |
之后我尝试用下面的代码给dm1添加标记列,但运行失败了:
dm1 = dm1.withColumn("highest_check", F.when(dm2.pk==dm1.pk & (dm2.number==dm1.number), True).otherwise(False))
正确的解决方法
其实不用额外生成新的DataFrame,直接在原表上用窗口函数就能一步搞定,这种方法更高效也更简洁:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义按pk分组的窗口规则 window_spec = Window.partitionBy("pk") # 直接添加highest_check列,判断当前行的number是否等于分组内的最大值 dm_result = dm1.withColumn( "highest_check", F.col("number") == F.max("number").over(window_spec) ) # 查看结果 dm_result.show()
为什么之前的方法失败?
你之前的代码里直接在withColumn中引用另一个DataFrame的列,Spark无法识别这种跨DataFrame的列关联,必须通过join操作才能实现关联匹配。如果一定要沿用你最初的思路,正确的做法是用左关联:
# 先获取每个pk分组中最大值对应的行,同时标记为True dm2 = dm1.withColumn('maxnumber', F.max('number').over(w))\ .where(F.col('number') == F.col('maxnumber'))\ .select("pk", "number", F.lit(True).alias("highest_check")) # 左关联原表,没有匹配到的行就将highest_check设为False dm_result = dm1.join(dm2, on=["pk", "number"], how="left")\ .withColumn("highest_check", F.coalesce(F.col("highest_check"), F.lit(False)))
不过显然第一种方法(直接使用窗口函数)是最优解,不需要额外的关联操作,性能更优,代码也更简洁易读。
备注:内容来源于stack exchange,提问作者Hotpacalypse
相关产品推荐
相关产品推荐

