PySpark非常规Join需求:合并两DataFrame并补缺失值为0
解决方案:合并两个PySpark DataFrame并补全缺失值
核心思路
要实现两个DataFrame所有日期记录的合并,并将Dataframe1中缺失的Value补为0,关键步骤是:
- 提取两个DataFrame中共同的主键字段(
Date、Id、Owner),合并后去重得到所有需要保留的记录组合; - 将该主键组合表与
Dataframe1左连接,确保所有日期记录都被保留; - 使用
coalesce函数将缺失的Value值替换为0。
代码实现
1. 创建测试DataFrame
from pyspark.sql import SparkSession from pyspark.sql.functions import coalesce, lit # 初始化SparkSession spark = SparkSession.builder.appName("MergeTwoDFs").getOrCreate() # 构建Dataframe1 df1_data = [ ("20240101", 2, 1, 3, 100), ("20240110", 2, 1, 3, 200) ] df1 = spark.createDataFrame(df1_data, ["Date", "Id", "Owner", "Id2", "Value"]) # 构建Dataframe2 df2_data = [ ("20240111", 2, 1), ("20240112", 2, 1), ("20240103", 2, 1) ] df2 = spark.createDataFrame(df2_data, ["Date", "Id", "Owner"])
2. 合并主键并左连接补全数据
# 获取所有唯一的(Date, Id, Owner)组合 all_unique_records = df1.select("Date", "Id", "Owner") \ .union(df2.select("Date", "Id", "Owner")) \ .distinct() # 左连接Dataframe1,补全Value字段 result_df = all_unique_records.join(df1, on=["Date", "Id", "Owner"], how="left") \ .withColumn("Value", coalesce(df1["Value"], lit(0))) # 按日期排序查看结果 result_df.orderBy("Date").show()
3. 可选:补全Id2字段(如果需要)
如果Id2是与Id、Owner关联的固定值,可以先提取映射关系再补全:
# 获取Id-Owner对应的Id2映射 id2_mapping = df1.select("Id", "Owner", "Id2").distinct() # 先补全Id2,再关联Value result_df_with_id2 = all_unique_records.join(id2_mapping, on=["Id", "Owner"], how="left") \ .join(df1, on=["Date", "Id", "Owner", "Id2"], how="left") \ .withColumn("Value", coalesce(df1["Value"], lit(0))) result_df_with_id2.orderBy("Date").show()
输出结果
执行基础方案后的输出:
+--------+---+-----+----+-----+ | Date| Id|Owner| Id2|Value| +--------+---+-----+----+-----+ |20240101| 2| 1| 3| 100| |20240103| 2| 1|null| 0| |20240110| 2| 1| 3| 200| |20240111| 2| 1|null| 0| |20240112| 2| 1|null| 0| +--------+---+-----+----+-----+
执行补全Id2后的输出:
+--------+---+-----+----+-----+ | Date| Id|Owner| Id2|Value| +--------+---+-----+----+-----+ |20240101| 2| 1| 3| 100| |20240103| 2| 1| 3| 0| |20240110| 2| 1| 3| 200| |20240111| 2| 1| 3| 0| |20240112| 2| 1| 3| 0| +--------+---+-----+----+-----+
内容的提问来源于stack exchange,提问作者TalendDeveloper
相关产品推荐
相关产品推荐

