PySpark实现多数据集对应列逐行求和并生成新DataFrame
PySpark 多DataFrame按共同键关联并逐行求和实现方案
核心思路:先统一重命名各DataFrame的得分列避免冲突,再基于共同维度列关联所有数据集,最后对得分列逐行求和(处理空值)得到结果。
步骤1:重命名得分列
由于部分DataFrame的得分列名称重复(比如aa2和aa3都叫score),先将其重命名为唯一标识名称:
# 重命名各DataFrame的得分列 aa1_renamed = aa1.withColumnRenamed("scoreHrs", "score_activity") aa2_renamed = aa2.withColumnRenamed("score", "score_calories") aa3_renamed = aa3.withColumnRenamed("score", "score_who5")
步骤2:关联所有DataFrame
所有数据集共享维度列:UserID, DoctorID, Department, Company, MonthNumber, YearNumber, PeriodID,以此为键进行关联。可根据业务需求选择关联方式:
outer:保留所有存在的行,缺失的得分数据补0inner:仅保留所有数据集都匹配的行
手动关联(适合少量DataFrame)
from pyspark.sql import functions as f # 先关联aa1和aa2 joined_df = aa1_renamed.join( aa2_renamed, on=["UserID", "DoctorID", "Department", "Company", "MonthNumber", "YearNumber", "PeriodID"], how="outer" ) # 再关联aa3 joined_df = joined_df.join( aa3_renamed, on=["UserID", "DoctorID", "Department", "Company", "MonthNumber", "YearNumber", "PeriodID"], how="outer" )
通用批量关联(适合任意数量的DataFrame)
如果后续数据集数量超过3个,用functools.reduce实现批量关联更高效:
from functools import reduce # 将所有重命名后的DataFrame放入列表 df_list = [aa1_renamed, aa2_renamed, aa3_renamed] # 定义关联键 join_keys = ["UserID", "DoctorID", "Department", "Company", "MonthNumber", "YearNumber", "PeriodID"] # 批量关联所有DataFrame joined_df = reduce( lambda df1, df2: df1.join(df2, on=join_keys, how="outer"), df_list )
步骤3:计算逐行总和
关联后可能存在空值(部分维度组合仅在部分数据集中存在),用f.coalesce将空值替换为0后再求和:
# 定义所有得分列的名称 score_columns = ["score_activity", "score_calories", "score_who5"] # 计算总分,空值替换为0 result_df = joined_df.withColumn( "total_score", sum(f.coalesce(f.col(col), f.lit(0)) for col in score_columns) ) # 可选:仅保留维度列和总分列 result_df = result_df.select(join_keys + ["total_score"])
完整代码示例
from pyspark.sql import functions as f from functools import reduce # 1. 重命名得分列 aa1_renamed = aa1.withColumnRenamed("scoreHrs", "score_activity") aa2_renamed = aa2.withColumnRenamed("score", "score_calories") aa3_renamed = aa3.withColumnRenamed("score", "score_who5") # 2. 批量关联所有DataFrame df_list = [aa1_renamed, aa2_renamed, aa3_renamed] join_keys = ["UserID", "DoctorID", "Department", "Company", "MonthNumber", "YearNumber", "PeriodID"] joined_df = reduce(lambda df1, df2: df1.join(df2, on=join_keys, how="outer"), df_list) # 3. 计算总分 score_columns = ["score_activity", "score_calories", "score_who5"] result_df = joined_df.withColumn( "total_score", sum(f.coalesce(f.col(col), f.lit(0)) for col in score_columns) ).select(join_keys + ["total_score"]) # 查看结果 result_df.show()
扩展说明
如果后续新增数据集,只需:
- 对新数据集的得分列进行重命名
- 将其加入
df_list - 在
score_columns中添加对应的新得分列名称即可
内容的提问来源于stack exchange,提问作者gbarel
相关产品推荐
相关产品推荐

