Spark设过滤与分区裁剪仍扫全表触发分区数超限原因排查
问题描述
扫描表table时触发分区数超限错误,所有查询失败,错误信息如下:
Caused by: org.apache.hadoop.hive.ql.metadata.HiveException: MetaException(message:Number of partitions scanned (=1000) on table 'table' exceeds limit (=100). This is controlled on the metastore server by metastore.limit.partition.request.)
执行的Spark代码:
spark.table("table"). filter($"dt" === "2023-01-01"). show
spark-shell配置:
--conf spark.sql.hive.convertMetastoreOrc=false \ --conf spark.sql.hive.metastorePartitionPruning=true \
表定义:
CREATE EXTERNAL TABLE table( columns ... PARTITIONED BY (dt date) STORED AS ORC TBLPROPERTIES ('external.table.purge'='true', 'orc.compress'='ZLIB')
已设置dt字段过滤条件及上述配置,但Spark仍扫描全表,触发分区数限制。
原因分析
- 分区字段类型不匹配:表中
dt是date类型,但过滤条件用的是字符串"2023-01-01",Spark无法将该过滤条件正确推送到Hive Metastore,导致先拉取全部分区再做过滤,触发分区数上限。 - 分区过滤推送失效:虽然配置了
spark.sql.hive.metastorePartitionPruning=true,但类型不匹配会导致Spark无法生成Metastore可识别的分区过滤请求,Metastore只能返回全部分区列表。
解决方案
- 修正过滤条件的类型匹配:将过滤条件中的字符串转换为
date类型,确保过滤逻辑能推送到Metastore:// 方法1:用to_date转换字符串 spark.table("table") .filter($"dt" === to_date(lit("2023-01-01"))) .show() // 方法2:直接使用date字面量 spark.table("table") .filter($"dt" === date"2023-01-01") .show() - 验证过滤推送效果:执行
explain()查看Spark执行计划,确认计划中显示仅扫描指定分区,而非全表。 - 临时调整Metastore配置(可选):若上述方法无法立即生效且有权限修改Metastore,可临时增大
metastore.limit.partition.request的值,但这只是临时方案,核心还是要确保分区过滤条件能正确推送。
内容的提问来源于stack exchange,提问作者gjin
相关产品推荐
相关产品推荐

