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

Spark三种日期过滤查询方案:性能对比与基准测试咨询

Spark三种日期过滤方案的性能对比与基准测试指南

我设计了三种略有不同的Spark表查询方式,用于筛选目标日期范围的数据(主要查询单日数据)。查看它们的物理执行计划后发现存在差异,想明确以下两点:

  1. 内连接过滤、普通列过滤、分区列过滤三种方案中哪一种在计算效率和Spark内存占用上最优
  2. 如何在Spark中对此类方案进行基准测试(类似Python的timeit,但Spark迭代成本较高)

测试代码(含中文注释)

// 下面使用的`calendar`表只有一列`calendar_day`,包含从1970年1月1日到2050年12月31日的所有日期
// 为了公平对比,我们从calendar表中只筛选出一天的数据
val start_day = LocalDate.parse("2023-01-07", DateTimeFormatter.ofPattern("yyyy-MM-dd"))
val end_day = LocalDate.parse("2023-01-07", DateTimeFormatter.ofPattern("yyyy-MM-dd"))
val start_day_literal = lit(Timestamp.valueOf(start_day.atStartOfDay())).cast("timestamp")
val end_day_literal = lit(Timestamp.valueOf(end_day.atStartOfDay())).cast("timestamp")
val t_dates = calendar 
.filter(col("calendar_day") >= start_day_literal)
.filter(col("calendar_day") <= end_day_literal)
.select("calendar_day")
t_dates.show(10, false);
+-------------------+
|calendar_day       |
+-------------------+
|2023-01-07 00:00:00|
+-------------------+

// 测试方案一:使用广播内连接过滤数据
val test_one = spark
.table("my_data_table")
.withColumn("customer_id", col("beneficiary_external_id").cast("long"))
.withColumn("start_date_trunc", col("start_date").cast("date"))
.withColumn("end_date_nvl", expr("nvl(cast(end_date as date), to_timestamp('2050-12-31', 'yyyy-MM-dd'))"))
.join(
    broadcast(t_dates).as("td"),
    joinExprs = col("calendar_day").between(
        lowerBound = col("start_date_trunc"),
        upperBound = col("end_date_nvl")
    ),
    joinType = "inner"
)
.select(
    col("customer_id"),
    col("start_date"),
    col("start_date_trunc"),
    col("start_day"),
    col("end_date"),
    col("end_date_nvl")
)
//test_one.show(100, false)
test_one.explain()
test_one.count(); // 结果行数:1,999,005,179

== Physical Plan ==
AdaptiveSparkPlan(isFinalPlan=false)
+- Project [customer_id#163L, start_date#114, start_date_trunc#182, start_day#128, end_date#115, end_date_nvl#202]
   +- BroadcastNestedLoopJoin BuildRight, Inner, ((calendar_day#0 >= cast(start_date_trunc#182 as timestamp)) && (calendar_day#0 <= end_date_nvl#202))
      :- Project [start_date#114, end_date#115, start_day#128, cast(beneficiary_external_id#113 as bigint) AS customer_id#163L, cast(start_date#114 as date) AS start_date_trunc#182, coalesce(cast(cast(end_date#115 as date) as timestamp), 2556057600000000) AS end_date_nvl#202]
      :  +- Filter isnotnull(cast(start_date#114 as date))
      :     +- FileScan parquet my_data_table[beneficiary_external_id#113,start_date#114,end_date#115,start_day#128] Batched: true, Format: Parquet, Location: CatalogFileIndex[s3://aac161b3-4383-880d-29a9-bd3d8e0115d0/9496ff37-0aaf-4b35-a0bc-5363a9213205], PartitionFilters: [], PushedFilters: [], ReadSchema: struct<beneficiary_external_id:string,start_date:timestamp,end_date:timestamp>
      +- BroadcastExchange IdentityBroadcastMode
         +- Filter ((calendar_day#0 >= 1673049600000000) && (calendar_day#0 <= 1673049600000000))
            +- FileScan parquet [calendar_day#0] Batched: true, Format: Parquet, Location: InMemoryFileIndex[s3://my-bucket/calendar_days/], PartitionFilters: [], PushedFilters: [GreaterThanOrEqual(calendar_day,2023-01-07 00:00:00.0), LessThanOrEqual(calendar_day,2023-01-07 ..., ReadSchema: struct<calendar_day:timestamp>
test_one: org.apache.spark.sql.DataFrame = [customer_id: bigint, start_date: timestamp ... 4 more fields]
res4: Long = 1999005179



// 测试方案二:使用普通列直接过滤,替代内连接
val dataset_timestamp = Timestamp.valueOf("2023-01-07 00:00:00.0")

val test_two = spark
.table("my_data_table")
.withColumn("customer_id", col("beneficiary_external_id").cast("long"))
.withColumn("start_date_trunc", col("start_date").cast("date"))
.withColumn("end_date_nvl", expr("nvl(cast(end_date as date), to_timestamp('2050-12-31', 'yyyy-MM-dd'))"))
.filter(
    col("start_date_trunc") <= dataset_timestamp
    )
.filter(
    col("end_date_nvl") >= dataset_timestamp
    )
.select(
    col("customer_id"),
    col("start_date"),
    col("start_date_trunc"),
    col("start_day"),
    col("end_date"),
    col("end_date_nvl")
)
//test_two.show(100, false)
test_two.explain()
test_two.count(); // 结果行数:1,999,005,179(与方案一一致)

== Physical Plan ==
*(1) Project [cast(beneficiary_external_id#1 as bigint) AS customer_id#51L, start_date#2, cast(start_date#2 as date) AS start_date_trunc#70, start_day#16, end_date#3, coalesce(cast(cast(end_date#3 as date) as timestamp), 2556057600000000) AS end_date_nvl#90]
+- *(1) Filter ((isnotnull(start_date#2) && (cast(cast(start_date#2 as date) as timestamp) <= 1673049600000000)) && (coalesce(cast(cast(end_date#3 as date) as timestamp), 2556057600000000) >= 1673049600000000))
   +- *(1) FileScan parquet my_data_table[beneficiary_external_id#1,start_date#2,end_date#3,start_day#16] Batched: true, Format: Parquet, Location: CatalogFileIndex[s3://aac161b3-4383-880d-29a9-bd3d8e0115d0/9496ff37-0aaf-4b35-a0bc-5363a9213205], PartitionFilters: [], PushedFilters: [IsNotNull(start_date)], ReadSchema: struct<beneficiary_external_id:string,start_date:timestamp,end_date:timestamp>
dataset_timestamp: java.sql.Timestamp = 2023-01-07 00:00:00.0
test_two: org.apache.spark.sql.DataFrame = [customer_id: bigint, start_date: timestamp ... 4 more fields]
res1: Long = 1999005179



// 测试方案三:使用分区列`start_day`进行过滤(与方案二的唯一区别)
val dataset_timestamp = Timestamp.valueOf("2023-01-07 00:00:00.0")

val test_three = spark
.table("my_data_table")
.withColumn("customer_id", col("beneficiary_external_id").cast("long"))
.withColumn("start_date_trunc", col("start_date").cast("date"))
.withColumn("end_date_nvl", expr("nvl(cast(end_date as date), to_timestamp('2050-12-31', 'yyyy-MM-dd'))"))
.filter(
    col("start_day") <= dataset_timestamp // 此处使用分区列替代普通列
    )
.filter(
    col("end_date_nvl") >= dataset_timestamp
    )
.select(
    col("customer_id"),
    col("start_date"),
    col("start_date_trunc"),
    col("start_day"),
    col("end_date"),
    col("end_date_nvl")
)
//test_three.show(100, false)
test_three.explain()
test_three.count(); // 结果行数:1,999,005,179(与前两个方案一致)

== Physical Plan ==
*(1) Project [cast(beneficiary_external_id#130 as bigint) AS customer_id#163L, start_date#131, cast(start_date#131 as date) AS start_date_trunc#182, start_day#145, end_date#132, coalesce(cast(cast(end_date#132 as date) as timestamp), 2556057600000000) AS end_date_nvl#202]
+- *(1) Filter (coalesce(cast(cast(end_date#132 as date) as timestamp), 2556057600000000) >= 1673049600000000)
   +- *(1) FileScan parquet my_data_table[beneficiary_external_id#130,start_date#131,end_date#132,start_day#145] Batched: true, Format: Parquet, Location: PrunedInMemoryFileIndex[s3://aac161b3-4383-880d-29a9-bd3d8e0115d0/9496ff37-0aaf-4b35-a0bc-5363..., PartitionFilters: [(start_day#145 <= 1673049600000000)], PushedFilters: [], ReadSchema: struct<beneficiary_external_id:string,start_date:timestamp,end_date:timestamp>
dataset_timestamp: java.sql.Timestamp = 2023-01-07 00:00:00.0
test_three: org.apache.spark.sql.DataFrame = [customer_id: bigint, start_date: timestamp ... 4 more fields]
res2: Long = 1999005179

一、三种方案的性能对比

1. 方案三(分区列过滤)最优

从物理执行计划和Spark优化逻辑来看,方案三是性能最优的选择,核心原因:

  • 分区裁剪生效:物理计划显示PartitionFilters: [(start_day#145 <= 1673049600000000)],Spark直接跳过所有start_day大于目标日期的分区,大幅减少了需要读取的文件数量和数据量,这是Spark最有效的数据过滤手段。
  • 无额外计算开销:不需要执行join操作,避免了嵌套循环join带来的内存占用和计算延迟。
  • 过滤逻辑更轻量:直接使用分区列进行比较,不需要对start_date做类型转换,降低了单条数据的计算成本。

2. 方案二(普通列过滤)次之

方案二性能优于方案一,但不如方案三:

  • 没有join开销,但无法利用分区裁剪,需要全表扫描所有分区。
  • 过滤逻辑需要对start_date进行类型转换后再比较,计算成本高于方案三。

3. 方案一(内连接过滤)性能最差

方案一的性能是三者中最差的:

  • 需要执行BroadcastNestedLoopJoin,即使使用了广播表,嵌套循环join的时间复杂度为O(n*m),在数据量极大时会带来显著的计算开销。
  • 未利用分区裁剪,需要全表扫描所有数据,同时还要处理join逻辑,内存占用和计算时间明显高于另外两种方案。

二、Spark基准测试方法

由于Spark迭代成本高,不能直接复用Python timeit的高频迭代方式,推荐以下实用方法:

1. 单次测试+多次重复取平均值

每次测试前清除缓存,避免缓存干扰结果,重复3-5次后取平均值:

def benchmark(name: String, df: => DataFrame): Long = {
  spark.catalog.clearCache() // 清除所有缓存
  val start = System.currentTimeMillis()
  df.count() // 触发执行
  val end = System.currentTimeMillis()
  val duration = end - start
  println(s"$name 执行时间: $duration ms")
  duration
}

// 执行测试
val test1Time = benchmark("方案一", test_one)
val test2Time = benchmark("方案二", test_two)
val test3Time = benchmark("方案三", test_three)

2. 利用Spark UI分析细节

通过Spark UI的Stage页面对比核心指标:

  • 读取的数据量(Input Size):确认分区裁剪是否生效
  • 任务执行时间(Duration):对比整体耗时
  • 内存占用(Executor Memory Usage):查看内存消耗差异
  • 查看SQL页面的物理执行计划,验证优化逻辑是否生效。

3. 模拟生产环境配置

测试时必须使用与生产环境一致的Spark配置(executor数量、内存、CPU核心数等),同时使用相同规模的测试数据,避免小数据量测试结果失真。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 19:22:08