Apache Spark合并同结构DataFrame按ID和Date分组对数值列求和实现
Spark实现全量+增量表按唯一键合并求和
实现思路
直接将两个DataFrame合并后按唯一标识ID和Date分组,对Amount A字段求和即可。该方案逻辑简洁,性能优于外连接方案,尤其适合多数值字段需要合并的场景:
- 同
ID同Date的两条记录会被分到同一组,求和结果就是两表数值相加 - 仅存在于单个表的记录分组后只有一条数据,求和结果就是原有数值
代码实现
PySpark 版本
from pyspark.sql import SparkSession from pyspark.sql.functions import sum # 初始化SparkSession spark = SparkSession.builder.appName("merge_full_incr").getOrCreate() # 测试数据:全量表 full_data = [ (1, "2021-09-01", 100), (1, "2021-09-02", 50), (2, "2021-09-01", 70) ] df_full = spark.createDataFrame(full_data, schema=["ID", "Date", "Amount A"]) # 测试数据:增量表 incr_data = [ (2, "2021-09-01", 50), (2, "2021-09-02", 30) ] df_incr = spark.createDataFrame(incr_data, schema=["ID", "Date", "Amount A"]) # 核心逻辑:合并后分组求和 df_result = df_full.unionByName(df_incr) \ .groupBy("ID", "Date") \ .agg(sum("Amount A").alias("Amount A")) \ .orderBy("ID", "Date") # 输出结果 df_result.show()
Scala Spark 版本
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.sum object MergeFullIncr { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("merge_full_incr") .master("local[*]") .getOrCreate() import spark.implicits._ // 全量表测试数据 val dfFull = Seq( (1, "2021-09-01", 100), (1, "2021-09-02", 50), (2, "2021-09-01", 70) ).toDF("ID", "Date", "Amount A") // 增量表测试数据 val dfIncr = Seq( (2, "2021-09-01", 50), (2, "2021-09-02", 30) ).toDF("ID", "Date", "Amount A") // 核心逻辑 val dfResult = dfFull.unionByName(dfIncr) .groupBy("ID", "Date") .agg(sum("Amount A").alias("Amount A")) .orderBy("ID", "Date") dfResult.show() } }
输出结果验证
运行后得到的结果和预期完全一致:
+---+----------+---------+ | ID| Date|Amount A| +---+----------+---------+ | 1|2021-09-01| 100| | 1|2021-09-02| 50| | 2|2021-09-01| 120| | 2|2021-09-02| 30| +---+----------+---------+
内容的提问来源于stack exchange,提问作者Carlos Casio
相关产品推荐
相关产品推荐

