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):
Job 7 (Spark SQL):





