PySpark DataFrame关联后填充field1、field2列值的技术问询
PySpark DataFrame关联与字段替换需求
输入DataFrame 1
| id | field1 | field2 |
|---|---|---|
| 1 | 10 | 20 |
| 2 | 30 | 40 |
| 3 | 50 | 60 |
输入DataFrame 2
| id | field_name | field_value |
|---|---|---|
| 1 | field1 | 100 |
| 1 | field2 | 200 |
| 2 | field1 | 300 |
初始关联后的DataFrame
| id | field1 | field2 | field_name | field_value |
|---|---|---|---|---|
| 1 | 10 | 20 | field1 | 100 |
| 1 | 10 | 20 | field2 | 200 |
| 2 | 30 | 40 | field1 | 300 |
预期输出结果
| id | field1 | field2 | field_name | field_value |
|---|---|---|---|---|
| 1 | 100 | 20 | field1 | 100 |
| 1 | 10 | 200 | field2 | 200 |
| 2 | 300 | 40 | field1 | 300 |
实现代码
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()
说明
- 先通过
join按id关联两个DataFrame,使用inner连接匹配示例中的初始关联结果。 - 使用
when条件函数:当field_name等于目标字段名时,用field_value替换对应字段的原值,否则保留原字段值。
内容的提问来源于stack exchange,提问作者kanishk kashyap
相关产品推荐
相关产品推荐

