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

Scala/Spark 2.4动态切换双表数据源取近6个月数据方案咨询

Scala/Spark 2.4 双表取数实现方案

核心思路

先确定新表Table2已覆盖的月份数,再根据规则拆分老表Table1和新表Table2的取数范围,最后合并两张表的有效数据。


1. 获取Table2的数据覆盖月份

假设表中有日期字段event_date(日期类型),先读取Table2的最小和最大日期,计算它覆盖了多少个完整月份:

import org.apache.spark.sql.functions._
import java.time.LocalDate
import java.time.temporal.ChronoUnit

// 读取Table2并提取最小、最大日期
val table2DF = spark.table("Table2")
val dateBounds = table2DF.agg(min("event_date").alias("min_dt"), max("event_date").alias("max_dt")).head()
val table2MinDt = dateBounds.getAs[LocalDate]("min_dt")
val table2MaxDt = dateBounds.getAs[LocalDate]("max_dt")

// 计算Table2覆盖的月份数(比如2024-02到2024-04算3个月)
val table2MonthCount = ChronoUnit.MONTHS.between(
  table2MinDt.withDayOfMonth(1), 
  table2MaxDt.withDayOfMonth(1)
) + 1

2. 拆分两张表的取数时间范围

以当前日期为基准,计算最近6个月的时间区间,再根据Table2的月份数分配取数范围:

// 计算最近6个月的时间边界:起始为当前月往前推6个月的1号,结束为当前月最后一天
val currentDt = LocalDate.now()
val sixMonthsAgoStart = currentDt.minusMonths(6).withDayOfMonth(1)
val currentMonthEnd = currentDt.withDayOfMonth(currentDt.lengthOfMonth())

// Table2只取最近6个月内的数据
val filteredTable2 = table2DF.filter(col("event_date").between(sixMonthsAgoStart, currentMonthEnd))

// 初始化最终结果为Table2的过滤数据
var finalData = filteredTable2

// 当Table2不足6个月时,从Table1补全剩余月份的数据
if (table2MonthCount < 6) {
  // Table1取数范围:从最近6个月起始,到Table2最早月份的前一个月月底
  val table1EndDt = table2MinDt.minusMonths(1).withDayOfMonth(
    table2MinDt.minusMonths(1).lengthOfMonth()
  )
  val filteredTable1 = spark.table("Table1").filter(
    col("event_date").between(sixMonthsAgoStart, table1EndDt)
  )
  // 合并两张表数据(确保Schema一致)
  finalData = finalData.union(filteredTable1)
}

3. 关键细节处理

  • Schema对齐:如果两张表字段不一致,必须先通过select指定相同字段列表,否则union会报错。示例:
    val alignedTable1 = filteredTable1.select("event_date", "col1", "col2", ...)
    val alignedTable2 = filteredTable2.select("event_date", "col1", "col2", ...)
    finalData = alignedTable2.union(alignedTable1)
    
  • 字符串日期转换:如果日期字段是字符串格式(如yyyy-MM-dd),需要先转成日期类型:
    val table2DF = spark.table("Table2").withColumn("event_date", to_date(col("event_date"), "yyyy-MM-dd"))
    
  • 调试日志:可以加日志输出时间范围和月份数,方便排查问题:
    println(s"Table2覆盖月份数: $table2MonthCount")
    println(s"Table1取数范围: $sixMonthsAgoStart 至 $table1EndDt")
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 06:35:16