I am running below spark code, which writes into a hive partitioned table.
df.write.mode(SaveMode.Overwrite).format("orc").insertInto("s**000h.test")
Internally all the executors are writing into the Hive stage area(.hive-staging_hive_2020-03-30_13-47-16_727_5670185411499574661-1) and it is taking more time as compared to while I am writing data explicitly into an HDFS directory as given below.
df.write.mode(mode).format("orc").partitionBy("dept_id").save(tempPath)
The time difference is coming to be around 1 hr for 900 partitions.
Could you please explain this behavior.