You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 19:06:02