PySpark实现:为A_rk与B_rk相同的行生成Pair列
PySpark实现DataFrame匹配Pair列需求
需求说明
现有如下Spark DataFrame:
+---+----+----+ | Id|A_rk|B_rk| +---+----+----+ | a| 5| 4| | b| 7| 7| | c| 5| 4| | d| 1| 0| +---+----+----+
需要新增名为Pair的列,规则如下:
- 当某行的
A_rk和B_rk组合在DataFrame中出现至少两次时,Pair列取值为该行的B_rk - 若该组合仅出现一次(无匹配行),
Pair列取值为0
期望输出结果:
+---+----+----+----+ | Id|A_rk|B_rk|Pair| +---+----+----+----+ | a| 5| 4| 4| | b| 7| 7| 0| | c| 5| 4| 4| | d| 1| 0| 0| +---+----+----+----+
实现方案
通过Spark分布式计算的分组统计或窗口函数实现,性能远优于Pandas循环。
方案1:分组统计+关联回原表
先统计每个(A_rk, B_rk)组合的出现次数,再关联回原表生成目标列:
from pyspark.sql import functions as F # 统计每个组合的出现次数 count_df = df.groupBy("A_rk", "B_rk").agg(F.count("Id").alias("cnt")) # 关联原表并生成Pair列 result_df = df.join(count_df, on=["A_rk", "B_rk"], how="left") \ .withColumn("Pair", F.when(F.col("cnt") >= 2, F.col("B_rk")).otherwise(0)) \ .drop("cnt")
方案2:窗口函数直接计算
无需额外关联,用窗口函数在原表上直接统计组合出现次数并生成目标列:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按A_rk和B_rk分组 window_spec = Window.partitionBy("A_rk", "B_rk") # 计算次数并生成Pair列 result_df = df.withColumn("cnt", F.count("Id").over(window_spec)) \ .withColumn("Pair", F.when(F.col("cnt") >= 2, F.col("B_rk")).otherwise(0)) \ .drop("cnt")
验证结果
执行上述任意方案后,result_df即可得到符合要求的输出。
内容的提问来源于stack exchange,提问作者devCharaf
相关产品推荐
相关产品推荐

