Spark三种日期过滤查询方案:性能对比与基准测试咨询
Spark三种日期过滤方案的性能对比与基准测试指南
我设计了三种略有不同的Spark表查询方式,用于筛选目标日期范围的数据(主要查询单日数据)。查看它们的物理执行计划后发现存在差异,想明确以下两点:
- 内连接过滤、普通列过滤、分区列过滤三种方案中哪一种在计算效率和Spark内存占用上最优
- 如何在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
相关产品推荐
相关产品推荐

