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

PySpark DataFrame关联后填充field1、field2列值的技术问询

PySpark DataFrame关联与字段替换需求

输入DataFrame 1

idfield1field2
11020
23040
35060

输入DataFrame 2

idfield_namefield_value
1field1100
1field2200
2field1300

初始关联后的DataFrame

idfield1field2field_namefield_value
11020field1100
11020field2200
23040field1300

预期输出结果

idfield1field2field_namefield_value
110020field1100
110200field2200
230040field1300

实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import when

# 初始化Spark会话
spark = SparkSession.builder.appName("FieldReplace").getOrCreate()

# 构造输入DataFrame1
df1 = spark.createDataFrame(
    [(1, 10, 20), (2, 30, 40), (3, 50, 60)],
    schema=["id", "field1", "field2"]
)

# 构造输入DataFrame2
df2 = spark.createDataFrame(
    [(1, "field1", 100), (1, "field2", 200), (2, "field1", 300)],
    schema=["id", "field_name", "field_value"]
)

# 关联两个DataFrame
joined_df = df1.join(df2, on="id", how="inner")

# 根据field_name替换对应字段的值
result_df = joined_df.withColumn(
    "field1",
    when(joined_df.field_name == "field1", joined_df.field_value).otherwise(joined_df.field1)
).withColumn(
    "field2",
    when(joined_df.field_name == "field2", joined_df.field_value).otherwise(joined_df.field2)
)

# 查看结果
result_df.show()

说明

  1. 先通过join按id关联两个DataFrame,使用inner连接匹配示例中的初始关联结果。
  2. 使用when条件函数:当field_name等于目标字段名时,用field_value替换对应字段的原值,否则保留原字段值。

内容的提问来源于stack exchange,提问作者kanishk kashyap

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 11:15:31