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

如何将行数不同的两个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()

代码说明

  1. 全外连接:how="full_outer"会保留两个DataFrame的所有行,匹配到的行合并列,未匹配的行对应字段自动填充Null。
  2. 关联键:通过Name、ID确保是同一主体的数据,再通过版本字段关联对应批次的特征,保证数据匹配准确性。
  3. 列选择:过滤掉需求中不需要的future_features_height,由于等值连接后Name和ID列完全一致,直接保留任意一侧即可。
  4. 排序:通过orderBy保证结果行顺序与示例一致。

内容的提问来源于stack exchange,提问作者Vivek Kaushik

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 04:35:32