Spark Filter Pushdown与数据加载交互机制:执行逻辑解惑
Spark Filter Pushdown与Partition Pruning工作原理解惑
核心前提:Spark的惰性执行模型
首先明确:Spark的Transformation操作(比如read.csv、filter)本身不会触发数据的实际读取或计算,它们只是构建一个逻辑执行计划。只有当你调用Action操作(如show()、count()、write())时,Spark才会将逻辑计划转换为物理计划,真正开始执行数据处理流程。
你看到的两个job,其中一个是inferSchema=True导致的:为了推断列的类型,Spark需要采样部分数据(默认读取少量文件或文件的前几行),这会触发一个轻量的job,但这不是全量加载数据到内存,只是采样动作。
Filter Pushdown的实际工作逻辑
以你的代码为例:
df = spark.read.csv('path', header=True, inferSchema=True) df_filtered = df.filter(col('salary') > 10000)
当你执行Action(比如df_filtered.show())时,Catalyst优化器会做以下事情:
- 将
filter操作下推到数据源读取阶段,避免全量加载数据后再过滤。 - 针对不同数据源,pushdown的实现细节有差异:
- Parquet/Orc等列式存储:这类格式自带元数据(比如列的统计信息),Spark可以直接利用这些元数据跳过不符合条件的列/行,甚至不需要解析整个文件,物理计划里的
PushedFilters会明确显示被下推的过滤条件(就是你看到的PushedFilters: [IsNotNull(country_code), EqualTo(country_code,IND)])。 - CSV等文本格式:虽然没有列存储的元数据优势,但Spark会在解析每一行CSV的过程中,直接判断是否符合filter条件,不符合的行不会被加载到内存,而是直接丢弃,本质上也是在数据源层面完成过滤,而非全量加载后再处理。
- Parquet/Orc等列式存储:这类格式自带元数据(比如列的统计信息),Spark可以直接利用这些元数据跳过不符合条件的列/行,甚至不需要解析整个文件,物理计划里的
你之前的误解在于:误以为spark.read.csv已经把数据加载到内存,但实际上read.csv只是定义了数据读取的逻辑,没有实际执行加载;直到Action触发后,才会结合filter条件,只读取符合要求的数据。
Partition Pruning的工作原理
Partition Pruning是针对分区表的优化(比如按country_code、date等字段分区),原理是:
- 当查询中包含对分区键的过滤条件时,Spark会直接跳过不符合条件的分区目录,根本不会读取这些分区下的文件。
- 比如你的物理计划里
PartitionFilters: [],说明这张表没有分区,所以没有触发分区修剪;如果是分区表,这里会显示被过滤的分区条件,物理计划中也会看到Spark只扫描符合条件的分区路径。
总结你的困惑点
spark.read.csv没有全量加载数据:它只是构建逻辑计划,只有Action触发才会执行实际读取。inferSchema触发的job是采样,不是全量加载:这个job只读取少量数据用于推断Schema,不影响后续的filter pushdown。- Filter Pushdown是在读取阶段执行:无论数据源是Parquet还是CSV,Catalyst都会将过滤逻辑下推到数据读取的最前端,只加载符合条件的记录,避免内存浪费。
内容的提问来源于stack exchange,提问作者DumbCoder
相关产品推荐
相关产品推荐

