Spark 2.4.4使用OR组合过滤条件时PartitionFilters与PushedFilters为空问题
问题原因
这个问题是Spark 2.4.4版本的Catalyst优化器规则限制导致的,具体原因如下:
- 首先,Spark 2.4.x版本没有实现跨OR分支的公共分区谓词提取优化规则:你给出的OR条件,人工可以逻辑推导简化为
date_col >= 20210928,但该版本的优化器不会自动对OR连接的多个分支做布尔逻辑合并,无法从混合了分区列、非分区列的OR条件里提取出公共的分区列过滤条件。 - 其次,OR两侧的分支都同时包含分区列
date_col和非分区列datetime_col的过滤条件,触发了下推阻断逻辑:Spark 2.4的谓词下推规则处理OR条件时,要求每个分支的过滤条件可以完全下推才会生效,只要分支中存在无法直接下推的混合判断,就会直接放弃整个条件的分区裁剪和谓词下推,因此物理计划中PartitionFilters和PushedFilters都为空,扫描全部分区。 - 该场景的下推优化在Spark 3.1及以上版本已经修复,新版本优化器可以自动提取OR分支中的公共分区过滤条件做分区裁剪。
临时解决方案
你可以手动简化过滤逻辑,先做分区裁剪再做业务条件过滤,就能让下推生效:
-- 先手动提取公共的分区列过滤条件做分区裁剪,减少扫描数据量 df.filter("date_col >= date_format(current_timestamp - interval 60 hours,\"yyyyMMdd\")") .filter("(datetime_col >= current_timestamp - interval 48 hours) or (datetime_col >= current_timestamp - interval 60 hours)")
内容的提问来源于stack exchange,提问作者Golokesh Patra
相关产品推荐
相关产品推荐

