PySpark分组后按优先级过滤DataFrame保留最高优先级status的方法
解题思路
- 第一步:映射优先级权重:按照给定的优先级规则,给三种status分别分配权重值,
OAOS-STP对应1,OAOS-nonSTP对应2,manual对应3,权重数值越小代表优先级越高。 - 第二步:新增辅助列:在原始DataFrame中新增优先级权重列,用于后续排序判断。
- 第三步:分组排序打标:用窗口函数按
unique-id分组,组内按优先级权重升序排序,给每行生成组内排名。 - 第四步:过滤结果:仅保留每个分组内排名为1的行,删除辅助列后即为目标结果。
PySpark实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import when, row_number from pyspark.sql.window import Window # 初始化SparkSession(已有实例可跳过该步骤) spark = SparkSession.builder.appName("priority_filter").getOrCreate() # 示例数据构造,实际使用时替换为你的源DataFrame读取逻辑即可 data = [ (1, "OAOS-STP"), (1, "OAOS-nonSTP"), (1, "manual"), (2, "OAOS-nonSTP"), (2, "manual"), (3, "OAOS-STP"), (3, "OAOS-nonSTP"), (4, "OAOS-STP"), (4, "manual") ] df = spark.createDataFrame(data, schema=["unique-id", "status"]) # 新增优先级权重列 df_with_priority = df.withColumn( "priority", when(df.status == "OAOS-STP", 1) .when(df.status == "OAOS-nonSTP", 2) .otherwise(3) ) # 定义窗口规则:按unique-id分区,按优先级升序排序 window_spec = Window.partitionBy("unique-id").orderBy("priority") # 生成组内排名 df_with_rank = df_with_priority.withColumn("rank", row_number().over(window_spec)) # 过滤取每个分组排名第一的行,删除辅助列得到最终结果 result_df = df_with_rank.filter(df_with_rank.rank == 1).drop("priority", "rank") # 输出查看结果 result_df.show()
内容的提问来源于stack exchange,提问作者vik
相关产品推荐
相关产品推荐

