Foundry中Spark如何启用分区裁剪(Partition Pruning)
问题背景
- 现有Palantir Foundry实例上维护一个增量构建模式的大型时序数据集,规模为565亿行、10个字段、总大小965GB,时间戳按1小时粒度分桶存储,每日新增数据量约10GB。
- 为优化查询性能,当前基于
measure_date和measuring_time两个字段做重分区,分区逻辑匹配业务访问模式:绝大多数查询都会按measure_date做过滤,measuring_time作为二级分区既可以减小Parquet文件体积,也适配常见的时间维度过滤场景。现有分区实现代码如下:
if ctx.is_incremental: return df.repartition(24, "measure_date", "measuring_time") else: return df.repartition(2200, "measure_date", "measuring_time")
- 哈希分区策略导致的文件大小不均衡问题不属于本次咨询处理范围。
- 现存异常:Foundry上的Spark引擎不会在带过滤条件的查询中触发分区裁剪。测试用SQL如下:
SELECT * FROM telemetry_data where measure_date = '2022-06-05'
- 异常佐证:查看该构建任务的物理查询计划,
PartitionFilters字段为空,Spark未识别利用现有分区,执行计划片段如下:
Batched: true, BucketedScan: false, DataFilters: [isnotnull(measure_date#170), (measure_date#170 = 19148)], Format: Parquet, Location: InMemoryFileIndex[sparkfoundry://prodapp06.palantir:8101/datasets/ri.foundry.main.dataset.xxx..., PartitionFilters: [], PushedFilters: [IsNotNull(measure_date), EqualTo(measure_date,2022-06-05)], ReadSchema: struct<xxx,measure_date:date,measuring_time_cet:timestamp,fxxx, ScanMode: RegularMode
根因说明
Spark的目录分区裁剪能力仅对Hive风格的分区存储结构生效,这类结构会把分区字段的键值直接嵌入到文件存储路径中,扫描时可以直接根据过滤条件匹配路径,跳过完全不需要访问的目录。
你当前使用的repartition()方法只会控制Spark写入阶段的内存数据分布、调整写入并行度,不会在存储层生成带分区字段的目录结构,Spark扫描时无法从路径层面直接跳过无关文件,自然不会触发PartitionFilters逻辑。当前执行计划中出现的PushedFilters只是Parquet文件级别的谓词下推,需要先遍历所有文件、读取Parquet文件脚注的统计信息才能做过滤,性能远低于目录级分区裁剪。
调整步骤
- 修改写入逻辑,声明存储层分区列
保留现有repartition()逻辑控制写入并行度、避免生成过多小文件,同时在写入时通过partitionBy声明分区字段,让存储层生成Hive风格的分区目录。如果是直接返回DataFrame给Foundry Transforms框架的写法,可以在输出数据集的配置中指定分区列,和调用partitionBy效果一致。
调整后的代码逻辑参考:if ctx.is_incremental: df = df.repartition(24, "measure_date", "measuring_time") else: df = df.repartition(2200, "measure_date", "measuring_time") # 核心:声明存储层分区字段,生成目录式分区结构 return df.partitionBy("measure_date", "measuring_time") - 全量重建数据集
分区结构属于数据集存储层面的元数据,增量构建无法修改历史数据的存储结构,需要触发一次全量构建,让所有历史数据按照新的分区规则组织目录。构建完成后可以在数据集的文件浏览页看到类似measure_date=2022-06-05/measuring_time=xxx/的目录层级,说明分区结构生效。 - 检查Spark配置项
确认构建任务的Spark配置没有手动关闭分区裁剪相关参数,Foundry环境下这些参数默认是开启的,无需额外修改:spark.sql.optimizer.dynamicPartitionPruning.enabled=true:开启动态分区裁剪spark.sql.hive.verifyPartitionPath=true:开启分区路径校验
- 验证裁剪效果
重新执行测试查询,查看物理执行计划,此时PartitionFilters字段会出现对应的measure_date过滤条件,扫描时只会命中符合条件的分区目录,扫描数据量会大幅下降。
注意:如果二级分区
measuring_time的基数过高(比如达到小时级、单表分区数超过10万),不建议把它设为二级存储分区,过多的分区目录会导致Driver端枚举分区路径的耗时大幅上涨,反而拖慢查询性能。这种场景下可以仅用measure_date做存储分区,measuring_time只做写入时的重分区字段控制文件大小即可。
内容的提问来源于stack exchange,提问作者twinkle2
相关产品推荐
相关产品推荐

