Spark如何优化分区数据读取?两种读取方式性能是否相当?
问题背景与疑问
数据存储结构
假设数据按以下分区结构存储:
data |---dt=2022-01-01 | |---app=x | |---app=y . . . |---dt=2022-01-02 | |---app=x | |---app=y . . . . . .
数据读取方式
若通过以下代码读取指定日期范围、指定app的Parquet数据:
where_condition = \ (col('dt').between(start_date, end_date)) & \ (col('app') == app_id) spark.spark_session.read.parquet(self.base_path).where(where_condition)
疑问
- 该方式与直接通过目标目录全路径列表读取数据的速度是否相当?
- Spark如何针对分区数据优化读取操作?
- 其利用分区特性的机制是什么?
解答
1. 速度对比
两者速度基本相当,甚至Spark的where过滤方式在部分场景下更灵活高效。因为Spark的分区发现机制会自动解析dt、app这类分区列的目录信息,执行where条件时会提前过滤掉不符合要求的目录,不会扫描这些目录下的文件,和直接指定目标目录列表的效果本质一致。
2. Spark对分区数据的读取优化
- 分区裁剪(Partition Pruning):读取阶段就根据过滤条件排除不需要的分区目录,避免扫描无关数据文件,大幅减少IO开销。
- 元数据缓存:Spark会缓存分区的元数据(比如分区列的取值、对应目录路径),后续重复读取时可直接复用缓存,无需重复遍历目录树获取元数据。
- 并行读取优化:针对符合条件的分区,Spark会将不同分区的文件分配到不同Task中并行读取,充分利用集群计算资源。
3. 利用分区特性的机制
Spark读取分区表时,会先递归扫描根目录下的所有子目录,解析目录名中的key=value格式(如dt=2022-01-01),将这些键值对识别为分区列及其取值,再把分区列信息添加到DataFrame的元数据中。
当执行where过滤条件时,Spark的Catalyst优化器会识别出过滤条件中的分区列,将过滤逻辑下推到读取阶段:
- 根据过滤条件计算出符合要求的分区取值组合(比如
dt在指定范围且app等于目标值); - 仅扫描这些取值对应的目录下的数据文件,完全跳过不符合条件的分区目录;
- 读取文件时,分区列的取值无需从Parquet文件中解析,直接从目录名中获取,进一步减少数据解析开销。
内容的提问来源于stack exchange,提问作者nirkov
相关产品推荐
相关产品推荐

