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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 17:45:28