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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 11:03:15