【发布时间】:2017-10-30 23:09:28
【问题描述】:
将分区过滤器应用于从具有超过 30,000 个分区的 Hive (v2.1.0) 表读取的 Spark (v2.0.2/2.1.1) DataFrames 时遇到问题。我想知道推荐的方法是什么,以及我做错了什么(如果有的话),因为当前的行为是导致大量性能和可靠性问题的根源。
为了启用修剪,我使用以下 Spark/Hive 属性:
--conf spark.sql.hive.metastorePartitionPruning=true
在 spark-shell 中运行查询时,我可以看到分区提取是通过调用 ThriftHiveMetastore.Iface.get_partitions 进行的,但在没有任何过滤的情况下会意外发生:
val myTable = spark.table("db.table")
val myTableData = myTable
.filter("local_date = '2017-09-01' or local_date = '2017-09-02'")
.cache
// The HMS call invoked is:
// #get_partitions('db', 'table', -1)
如果我使用更简单的过滤器,则会根据需要过滤分区:
val myTableData = myTable
.filter("local_date = '2017-09-01'")
.cache
// The HMS call invoked is:
// #get_partitions_by_filter(
// 'db', 'table',
// 'local_date = "2017-09-01"',
// -1
// )
如果我重写过滤器以使用范围运算符而不是简单地检查是否相等,过滤也可以正常工作:
val myTableData = myTable
.filter("local_date >= '2017-09-01' and local_date <= '2017-09-02'")
.cache
// The HMS call invoked is:
// #get_partitions_by_filter(
// 'db', 'table',
// 'local_date >= '2017-09-01' and local_date <= '2017-09-02'',
// -1
// )
在我们的例子中,从性能角度来看,这种行为是有问题的;正确过滤后,通话时间在 4 分钟左右,而在 1 秒左右。此外,每次查询时例行地将大量Partition 对象加载到堆上最终会导致 Metastore 服务出现内存问题。
似乎某些类型的过滤器构造的解析和解释存在错误,但是我无法在 Spark JIRA 中找到相关问题。是否有首选方法或特定 Spark 版本,其中过滤器适用于所有过滤器变体?或者在构建过滤器时我必须使用特定的形式(例如范围运算符)吗?如果是这样,此限制是否记录在任何地方?
【问题讨论】:
-
在 Spark JIRA SPARK-22247987654322@提出了一个票
标签: apache-spark hive spark-dataframe