PySpark按Type分组比较两列最大值生成status列问题求助
解决PySpark按分组判断状态的问题
你的代码报错是因为直接在withColumn中使用聚合函数max()时,未指定分组上下文,PySpark无法确定分组依据,因此抛出grouping expressions sequence is empty错误。要实现按Type分组判断状态的需求,需要使用窗口函数来处理分组内的聚合计算。
完整解决方案代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义按Type分组的窗口 window_spec = Window.partitionBy("Type") # 计算分组内的最大值并生成status列 df = df.withColumn("max_n1", F.max("n1").over(window_spec)) \ .withColumn("max_n2", F.max("n2").over(window_spec)) \ .withColumn("status", F.when(F.col("max_n1") > F.col("max_n2"), "OK").otherwise("KO")) \ .drop("max_n1", "max_n2") # 移除中间计算列(可选)
代码说明
- 窗口定义:
Window.partitionBy("Type")创建了以Type为分组依据的窗口,同一Type下的所有行会共享该分组的聚合结果。 - 分组最大值计算:通过
F.max().over(window_spec)分别计算每个Type分组内n1和n2的最大值,生成临时列max_n1和max_n2。 - 状态判断:使用
when...otherwise逻辑判断分组内max_n1是否大于max_n2,生成status列。 - 清理中间列:如果不需要临时的最大值列,用
drop()移除即可。
执行结果
运行代码后,会得到你期望的目标数据集:
| Type | n1 | n2 | status |
|---|---|---|---|
| A | 12 | 17 | OK |
| A | 14 | 16 | OK |
| B | 12 | 13 | KO |
| B | 12 | 11 | KO |
| A | 11 | 12 | OK |
| A | 19 | 15 | OK |
内容的提问来源于stack exchange,提问作者Nabs335
相关产品推荐
相关产品推荐

