PySpark实现两DataFrame合并并按用户ID汇总列值
PySpark实现类似Pandas的合并+列累加需求
嘿,我来帮你搞定这个问题!你之前尝试用union和agg没成功是因为union是用来堆叠同结构的行数据,而你要的是按user_id合并两个表的列,同时把Prev_table的earnings和New_table的profit相加,还要保留其他字段,正确的思路是用全外连接+聚合,下面一步步给你讲具体实现:
核心思路
- 先把
New_table的profit列重命名为earnings,统一字段名方便后续计算 - 对两个表做全外连接(full outer join),确保所有
user_id都被保留(不管只在哪个表出现) - 按
user_id分组聚合,计算earnings的总和,同时提取start_date和end_date的非空值 - 调整列顺序匹配期望结果
具体代码实现
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
相关产品推荐
相关产品推荐

