Spark Scala实现对齐有效月份的年度数据同比计算方法
Spark Scala 实现带月份对齐规则的年度销量同比计算
核心逻辑
计算相邻年度销量总和时,仅对两年共同存在的月份的销量做求和,单边缺失的月份数值直接剔除,保证两年统计口径完全对齐。
具体实现代码
import org.apache.spark.sql.functions._ // 1. 预处理:拆分年、月维度 val processedSalesDF = salesDF .withColumn("Year", substring(col("Month"), 1, 4).cast("Int")) .withColumn("MonthNum", substring(col("Month"), 6, 2).cast("Int")) .select("Year", "MonthNum", "Fruit", "Qty") // 2. 聚合得到每个水果、每个年度的存在月份集合 val yearMonthsDF = processedSalesDF .groupBy("Fruit", "Year") .agg(collect_set("MonthNum").alias("ExistMonths")) // 3. 关联相邻年度,计算对齐用的共有月份集合 val alignRuleDF = yearMonthsDF.alias("curr") .join( yearMonthsDF.alias("prev"), expr("curr.Fruit = prev.Fruit AND curr.Year = prev.Year + 1"), "left" ) .select( col("curr.Fruit"), col("curr.Year").alias("CurrYear"), col("prev.Year").alias("PrevYear"), coalesce( array_intersect(col("curr.ExistMonths"), col("prev.ExistMonths")), col("curr.ExistMonths") ).alias("AlignedMonths") ) // 4. 关联明细计算当年对齐后销量总和 val currSumDF = alignRuleDF .join( processedSalesDF, expr(""" alignRuleDF.Fruit = processedSalesDF.Fruit AND processedSalesDF.Year = alignRuleDF.CurrYear AND array_contains(alignRuleDF.AlignedMonths, processedSalesDF.MonthNum) """) ) .groupBy(col("CurrYear").alias("Year"), col("Fruit")) .agg(sum(col("Qty")).alias("sum_of_qty_This_year")) // 5. 关联明细计算前一年对齐后销量总和 val prevSumDF = alignRuleDF .where(col("PrevYear").isNotNull) .join( processedSalesDF, expr(""" alignRuleDF.Fruit = processedSalesDF.Fruit AND processedSalesDF.Year = alignRuleDF.PrevYear AND array_contains(alignRuleDF.AlignedMonths, processedSalesDF.MonthNum) """) ) .groupBy(col("CurrYear").alias("Year"), col("Fruit")) .agg(sum(col("Qty")).alias("sum_of_qty_previous_year")) // 6. 拼接得到最终结果 val resultDF = currSumDF .join(prevSumDF, Seq("Year", "Fruit"), "left") .orderBy("Fruit", "Year")
说明
- 所有相邻年度的对齐规则独立计算,比如2021-2022的共有月份、2020-2021的共有月份互不干扰,符合跨年度两两对比的要求
- 若使用Spark 2.4以下版本,可自行实现UDF替换
array_intersect,输入两个数组返回交集即可 - 数据集中最早的年度无更早对比数据时,
sum_of_qty_previous_year字段返回null,可根据业务需要用fill函数填充默认值
内容的提问来源于stack exchange,提问作者Gaurav
相关产品推荐
相关产品推荐

