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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 06:20:46