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

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:保留所有存在的行,缺失的得分数据补0
  • inner:仅保留所有数据集都匹配的行

手动关联(适合少量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()

扩展说明

如果后续新增数据集,只需:

  1. 对新数据集的得分列进行重命名
  2. 将其加入df_list
  3. 在score_columns中添加对应的新得分列名称即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 23:43:15