如何将Scala列表元素作为参数传入Spark filter表达式统计数据量
实现方案
你可以通过以下两种方式实现需求,推荐优先使用第二种性能更优的方案:
方案1:逐个遍历日期统计
适合日期数量少、表数据量小的场景,逻辑简单直接:
import scala.collection.mutable.ListBuffer import org.apache.spark.sql.functions.col // 你已有的日期列表 val dateList: ListBuffer[String] = ListBuffer("2021-10-01", "2021-10-02", "2021-10-03", "2021-10-04") // 遍历每个日期过滤统计 val dateCountResult = dateList.map { date => val cnt = spark.read.table(existingTable) .filter(col("event_date") === date) .count() (date, cnt) } // 输出结果 dateCountResult.foreach { case (dt, num) => println(s"$dt: $num 条") }
方案2:一次性分组统计(推荐)
只扫描一次表即可完成所有日期的统计,性能远高于方案1,适合大表、多日期统计场景:
import scala.collection.mutable.ListBuffer import org.apache.spark.sql.functions.col val dateList: ListBuffer[String] = ListBuffer("2021-10-01", "2021-10-02", "2021-10-03", "2021-10-04") val targetDates = dateList.toSeq // 先过滤出所有目标日期的数据,再分组统计 val countDF = spark.read.table(existingTable) .filter(col("event_date").isin(targetDates: _*)) .groupBy("event_date") .count() // 转成Map格式,补全无数据的日期(统计值为0) val countMap = countDF.collect() .map(row => row.getAs[String]("event_date") -> row.getAs[Long]("count")) .toMap val finalResult = dateList.map(dt => dt -> countMap.getOrElse(dt, 0L)).toMap // 输出结果 finalResult.foreach { case (dt, num) => println(s"$dt: $num 条") }
注意事项
- 如果表中
event_date字段是DateType类型而非字符串类型,需要把日期字符串转成日期类型再匹配,避免类型不匹配导致过滤失败:col("event_date") === to_date(lit(date), "yyyy-MM-dd") - 如果你的表是按
event_date做分区的分区表,两种方案都会自动触发分区裁剪,不会扫描其他日期的分区,性能会进一步提升。
内容的提问来源于stack exchange,提问作者SanjanaSanju
相关产品推荐
相关产品推荐

