what is spark doing after insertInto?

Viewed 78

I have a super-simple pyspark script:

  • Run query(Hive) and create a dataframe A
  • Perform aggregates on A which creates dataframe B
  • Print the number of rows on the aggregates with B.count()
  • Save the results of B in a Hive table using B.insertInto()

But I noticed something, in the Spark web UI the insertInto is listed as completed, but the client program(notebook) is still marking the insert as as running, if I run a count directly to the Hive table with a Hive client(no spark) the row-count is not the same as printed with B.count(), if I run the query again, the number of rows increases(but still not matching to B.count()) after some minutes, the row count hive query, matches B.count(). Question is, if the insertInto() job is already completed (according to the Spark web UI) what is it doing? given by the row-count increase behavior it seems as it is still running the insertInto but that does not matches with the spark web UI. My guess is something like hive table partition metadata update is running, or something similar.

0 Answers
Related