PySpark实现:基于匹配ID用DataFrame B更新DataFrame A的Scores列
在PySpark中根据ID匹配更新DataFrame列值
需求说明
将dfA中每个ID对应的所有Scores值,替换为dfB中该ID对应的Scores值,最终实现同一ID的Scores统一为dfB中的对应值。
示例数据准备
先创建示例中的两个DataFrame,方便复现验证:
from pyspark.sql import SparkSession from pyspark.sql.functions import col # 初始化Spark会话 spark = SparkSession.builder.appName("UpdateScoresDemo").getOrCreate() # 构建dfA data_a = [("A", 20), ("A", 40), ("A", 60), ("B", 10), ("B", 90)] dfA = spark.createDataFrame(data_a, schema=["ID", "Scores"]) # 构建dfB data_b = [("A", 60), ("B", 90)] dfB = spark.createDataFrame(data_b, schema=["ID", "Scores"])
解决方案
核心逻辑是通过ID关联两个DataFrame,用dfB的Scores覆盖dfA的对应列。由于dfB中每个ID仅存一条目标值记录,可通过join操作快速实现:
方法1:直接关联重命名列
# 用左连接保留dfA所有行(若ID不匹配则Scores为null) updated_df = dfA.join(dfB, on="ID", how="left") \ .select(col("ID"), col("Scores").alias("Scores")) # 查看结果 updated_df.show()
注:join后dfB的
Scores会默认覆盖同名列,也可以手动指定列别名避免混淆:
# 先给dfB的Scores重命名 dfB_alias = dfB.withColumnRenamed("Scores", "Target_Scores") # 关联后替换列 updated_df = dfA.join(dfB_alias, on="ID", how="left") \ .select("ID", "Target_Scores") \ .withColumnRenamed("Target_Scores", "Scores")
方法2:广播小表优化性能
如果dfB数据量很小,推荐使用broadcast广播小表,大幅提升关联性能:
from pyspark.sql.functions import broadcast updated_df = dfA.join(broadcast(dfB), on="ID", how="left") \ .select("ID", col("Scores").alias("Scores"))
输出结果
执行上述代码后,会得到期望的结果:
+---+------+ | ID|Scores| +---+------+ | A| 60| | A| 60| | A| 60| | B| 90| | B| 90| +---+------+
补充说明
- 若只需要保留dfB中存在的ID记录,将
how="left"改为how="inner"即可; - 若dfA存在dfB没有的ID,左连接会保留这些行,对应的
Scores会显示为null,可根据需求用fillna填充默认值。
内容的提问来源于stack exchange,提问作者Luis Valdes
相关产品推荐
相关产品推荐

