Based on same physical plan, why the spark sql is much faster than dataframe when i using WHERE to filter data from hive partitioned table?

Viewed 76

Describes

I have a hive partitioned table (table1 for example) which has 4 partitioned columns, and when I using spark-shell to filter the table with the first two partitioned columns(project and date), I got different performance. Using spark sql to execute sql is much faster than spark dataframe's WHERE, and what's even weirder is that they have same physical plan.

And I have checked the Spark Web UI, showing that the input and output recoders and data size for those two job(one is spark sql, one is dataframe) are same too.

So, it confuses me so much. I think that with the same physical plan, it should have same performance.And hope your help, thanks so much.

Here are some details:

hive version: 1.1.0-cdh5.8.5
spark version: 2.1.0.cloudera2
scala version: 2.11.8
jdk version: 1.8.0_181
hive table partitioned columns: project,date,type,section
hive table file format: SequenceFile

This is the code I ran on spark-shell base on yarn:

scala> spark.sql("select * from table1 where project='a' and `date`='b'").explain
== Physical Plan ==
HiveTableScan [...other table columns...., project#176, date#177, type#178, section#179], MetastoreRelation table1, [isnotnull(project#176), isnotnull(date#177), (project#176 = a), (date#177 = b)]

scala> spark.table("table1").where("project='a'").where("`date`='b'").explain
== Physical Plan ==
HiveTableScan [...other table columns..., project#195, date#196, type#197, section#198], MetastoreRelation table1, [isnotnull(project#195), (project#195 = a), isnotnull(date#196), (date#196 = b)]


// I have run those statements several times before, so it shouldn't be the cache problem
scala> spark.time(spark.sql("select * from table1 where project='a' and `date`='b'").count)
Time taken: 8146 ms                                                             
res10: Long = 45392

scala> spark.time(spark.table("table1").where("project='a'").where("`date`='b'").count)
Time taken: 513 ms
res11: Long = 45392

Spark Web UI snapshot:

Job 6(Dataframe):

stage info for Job6

stage 11 detail for Job 6

stage 12 detail for Job 6

Job 7 (Spark SQL):

stage info for Job7

stage 13 detail for Job 7

stage 14 detail for Job 7

0 Answers
Related