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

Spark中如何遍历Scala日期ListBuffer逐次读取数据避免数据量过大

Scala逐日期处理数据实现方案

是的,可以直接用Scala集合的foreach方法遍历ListBuffer完成单日期逐次处理,无需手动维护while循环索引,实现如下:

改造后完整代码

import org.apache.spark.sql.functions.col
import scala.collection.mutable.ListBuffer

// 原始日期列表
val partList: ListBuffer[String] = ListBuffer("2021-10-01", "2021-10-02", "2021-10-03", "2021-10-04", "2021-10-05", "2021-10-06", "2021-10-07", "2021-10-08")

// 遍历每个日期单独处理
partList.foreach { currentDate =>
  // 单次仅加载当前日期的源表数据
  val fctDF = ss.read.table(existingTable).filter(s"event_date = '${currentDate}'")
  
  // 当前日期无数据时跳过后续逻辑
  if (fctDF.count() > 0) {
    fctDF.createOrReplaceTempView("vw_exist_fct")
    val existingRecordsQuery = getExistingRecordsMergeQuery(azUpdateTS, key)
    ss.sql(existingRecordsQuery)
      .drop("az_insert_ts", "az_update_ts")
      .withColumn("az_insert_ts", col("new_az_insert_ts"))
      .withColumn("az_update_ts", col("new_az_update_ts"))
      .drop("new_az_insert_ts", "new_az_update_ts")
      .select(mrg_tbl_cols.head, mrg_tbl_cols.tail: _*)
      .coalesce(72 * 2)
      .write.mode("Append")
      .format("delta")
      .insertInto(mergeTable)
  }
}

// 所有日期处理完成后,统一写入最终目标表
val mergedDataDF = ss.read.table(mergeTable).coalesce(72 * 2)
mergedDataDF.coalesce(72)
  .write.mode("Overwrite")
  .format("delta")
  .insertInto(s"${tgtSchema}.${tgtTbl}")

核心调整说明

  • 筛选逻辑从全量in匹配修改为单日期=匹配,大幅降低单次加载的数据量,避免内存溢出
  • 最终目标表写入逻辑移到foreach外部,所有日期处理完成后统一生成,避免重复覆盖浪费计算资源
  • 用mrg_tbl_cols.head+mrg_tbl_cols.tail替代原来的索引切片写法,更符合Scala代码规范

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 19:45:06