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
相关产品推荐
相关产品推荐

