如何将行数不同的两个PySpark DataFrame合并为新增列形式?
合并行数不同的PySpark DataFrame(新增列形式)
需求说明
需要将两个行数可能不同的PySpark DataFrame合并,将其中一个的列作为新增列整合到另一个中。匹配逻辑基于Name、ID字段,同时关联对应的版本标识(past_features_version与future_features_version),不匹配的行对应字段填充Null。
示例数据
df1(历史特征数据)
Schema:
Name|ID|past_features_height|past_features_width|past_features_version
数据:
A|123|34|42|version1 A|123|30|45|version2 A|123|33|47|version3
df2(未来特征数据)
Schema:
Name|ID|future_features_height|future_features_width|future_features_version
数据:
A|123|32|45|version1 A|123|35|42|version2 A|123|37|43|version3 A|123|39|46|version4 A|123|40|32|version5
期望结果
Schema:
Name|ID|past_features_height|past_features_width|past_features_version|future_features_width|future_features_version
数据:
A|123|34|42|version1|32|45|version1 A|123|30|45|version2|35|42|version2 A|123|33|47|version3|37|43|version3 A|123|Null|Null|Null|39|46|version4 A|123|Null|Null|Null|40|32|version5
解决方案
使用PySpark的**全外连接(full outer join)**实现,指定关联键为Name、ID以及版本字段,最后选择需要的列(过滤掉不需要的future_features_height)。
代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col # 初始化SparkSession spark = SparkSession.builder.appName("DataFrameMerge").getOrCreate() # 创建df1数据 data1 = [ ("A", 123, 34, 42, "version1"), ("A", 123, 30, 45, "version2"), ("A", 123, 33, 47, "version3") ] schema1 = ["Name", "ID", "past_features_height", "past_features_width", "past_features_version"] df1 = spark.createDataFrame(data1, schema=schema1) # 创建df2数据 data2 = [ ("A", 123, 32, 45, "version1"), ("A", 123, 35, 42, "version2"), ("A", 123, 37, 43, "version3"), ("A", 123, 39, 46, "version4"), ("A", 123, 40, 32, "version5") ] schema2 = ["Name", "ID", "future_features_height", "future_features_width", "future_features_version"] df2 = spark.createDataFrame(data2, schema=schema2) # 执行全外连接 merged_df = df1.join( df2, (df1.Name == df2.Name) & (df1.ID == df2.ID) & (df1.past_features_version == df2.future_features_version), how="full_outer" ) # 选择需要的列,处理重复的Name和ID列 result_df = merged_df.select( col("Name"), col("ID"), "past_features_height", "past_features_width", "past_features_version", "future_features_width", "future_features_version" ).orderBy("future_features_version") # 按版本排序对齐示例结果 # 展示结果 result_df.show()
代码说明
- 全外连接:
how="full_outer"会保留两个DataFrame的所有行,匹配到的行合并列,未匹配的行对应字段自动填充Null。 - 关联键:通过
Name、ID确保是同一主体的数据,再通过版本字段关联对应批次的特征,保证数据匹配准确性。 - 列选择:过滤掉需求中不需要的
future_features_height,由于等值连接后Name和ID列完全一致,直接保留任意一侧即可。 - 排序:通过
orderBy保证结果行顺序与示例一致。
内容的提问来源于stack exchange,提问作者Vivek Kaushik
相关产品推荐
相关产品推荐

