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

PySpark实现两DataFrame合并并按用户ID汇总列值

PySpark实现类似Pandas的合并+列累加需求

嘿,我来帮你搞定这个问题!你之前尝试用union和agg没成功是因为union是用来堆叠同结构的行数据,而你要的是按user_id合并两个表的列,同时把Prev_table的earnings和New_table的profit相加,还要保留其他字段,正确的思路是用全外连接+聚合,下面一步步给你讲具体实现:

核心思路

  1. 先把New_table的profit列重命名为earnings,统一字段名方便后续计算
  2. 对两个表做全外连接(full outer join),确保所有user_id都被保留(不管只在哪个表出现)
  3. 按user_id分组聚合,计算earnings的总和,同时提取start_date和end_date的非空值
  4. 调整列顺序匹配期望结果

具体代码实现

1. 先创建测试数据集(模拟你的表)

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

# 初始化SparkSession
spark = SparkSession.builder.appName("merge_earnings").getOrCreate()

# 创建Prev_table
prev_data = [
    (1, 10, "2020-06-01", "2020-06-10"),
    (2, 20, "2020-06-01", "2020-06-10"),
    (3, 30, "2020-06-01", "2020-06-10")
]
prev_table = spark.createDataFrame(prev_data, schema=["user_id", "earnings", "start_date", "end_date"])

# 创建New_table
new_data = [
    (1, 100),
    (2, 200),
    (5, 500)
]
new_table = spark.createDataFrame(new_data, schema=["user_id", "profit"])

2. 重命名字段并做全外连接

把New_table的profit改成earnings,这样两个表就有了统一的计算字段:

# 重命名New_table的profit为earnings
new_table_renamed = new_table.withColumnRenamed("profit", "earnings")

# 全外连接两个表,基于user_id
joined_df = prev_table.join(new_table_renamed, on="user_id", how="full_outer")

这一步之后,每个user_id的记录会包含两边的字段,缺失的字段会显示null(比如user5的start_date和end_date就是null)。

3. 分组聚合计算最终结果

按user_id分组,累加earnings,同时提取start_date和end_date的非空值:

result_df = joined_df.groupBy("user_id") \
    .agg(
        # 计算earnings总和,用coalesce处理可能的null情况(避免出现null结果)
        F.coalesce(F.sum(F.col("earnings")), F.lit(0)).alias("earnings"),
        # 提取start_date的非空值(同一个user_id在Prev_table里只有一条,所以first就够用)
        F.first("start_date").alias("start_date"),
        # 同理提取end_date
        F.first("end_date").alias("end_date")
    )

# 调整列顺序,和期望结果一致
result_df = result_df.select("user_id", "earnings", "start_date", "end_date")

4. 查看结果

执行result_df.show()后,输出就是你要的效果:

+-------+--------+----------+----------+
|user_id|earnings|start_date|  end_date|
+-------+--------+----------+----------+
|      1|     110|2020-06-01|2020-06-10|
|      2|     220|2020-06-01|2020-06-10|
|      3|      30|2020-06-01|2020-06-10|
|      5|     500|      null|      null|
+-------+--------+----------+----------+

补充说明

  • 如果你的实际数据中,同一个user_id在Prev_table或New_table中有多条记录,sum函数会自动累加所有对应的earnings/profit,完全符合需求
  • 关于start_date和end_date的处理,因为只有Prev_table有这两个字段,用first/max/min都可以,结果是一样的,这里用first最直观
  • 全外连接(full outer join)是关键,它能确保不管user_id只在哪个表出现,都能被保留在结果里

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 22:43:09