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

