Is spark.read or spark.sql lazy transformations?

Viewed 593

In Spark if the source data has changed in between two action calls why I still get previous o/p not the most recent ones. Through DAG all operations will get executed including read operation once action is called. Isn't it?

e.g. df = spark.sql("select * from dummy.table1") #Reading from spark table which has two records into dataframe.

df.count() #Gives count as 2 records

Now, a record inserted into table and action is called withou re-running command1 .

df.count() #Still gives count as 2 records.

I was expecting Spark will execute read operation again and fetch total 3 records into dataframe.

Where my understanding is wrong ?

2 Answers

This is most related to inmutability as I see it. DataFrames are inmutables, hence changes in the original table are not reflected on them.

Once a dataframe is evaluated, it will be never calculated again. So once the dataframe named df is evaluated, it is the picture of table1 at the time of evaluation, it doesn't matter if table1 changes, df won't. So the second df.count does not trigger evaluation it just return the previous result, which is 2

If you want the desired results you have to load again the DF in a different variable:

val df = spark.sql("select * from dummy.table1")
df.count() //Will trigger evaluation and return 2

//Insert record

val df2 = spark.sql("select * from dummy.table1")
df2.count() //Will trigger evaluation and return 3

Or using var instead of val (which is bad)

var df = spark.sql("select * from dummy.table1")
df.count() //Will trigger evaluation and return 2

//Insert record

df = spark.sql("select * from dummy.table1")
df.count() //Will trigger evaluation and return 3

This said: yes, spark read and spark sql are lazy, those are not called until an action is found, but once that happens, evaluation won't be trigger ever again in that dataframe

To contrast your assertion, this below does give a difference - using Databricks Notebook (cells). The insert operation is not known that you indicate.

But the following using parquet or csv based Spark - thus not Hive table, does force a difference in results as the files making up the table change. For a DAG re-compute, the same set of files are used afaik, though.

//1st time in a cell
val df = spark.read.csv("/FileStore/tables/count.txt")
df.write.mode("append").saveAsTable("tab2")

//1st time in another cell
val df2 = spark.sql("select * from tab2")
df2.count() 
//4 is returned


//2nd time in a different cell
val df = spark.read.csv("/FileStore/tables/count.txt")
df.write.mode("append").saveAsTable("tab2")

//2nd time in another cell
df2.count() 
//8 is returned

Refutes your assertion. Also tried with .enableHiveSupport(), no difference.

Even if creating a Hive table directly in Databricks:

spark.sql("CREATE TABLE tab5 (id INT, name STRING, age INT) STORED AS ORC;")
spark.sql(""" INSERT INTO tab5 VALUES (1, 'Amy Smith', 7) """)

...
df.count()
...

spark.sql(""" INSERT INTO tab5 VALUES (2, 'Amy SmithS', 77) """)
df.count()

...

Still get updated counts.

However, the for a Hive created ORC Serde table, the following "hive" approach or using an insert via spark.sql:

val dfX = Seq((88,"John", 888)).toDF("id" ,"name", "age")
dfX.write.format("hive").mode("append").saveAsTable("tab5")

or

spark.sql(""" INSERT INTO tab5 VALUES (1, 'Amy Smith', 7) """)

will sometimes show and sometimes not show an updated count when just the 2nd df.count() is issued. This is due to Hive / Spark lack of synchronization that may depend on some internal flagging of changes. In any event not consistent. Double-checked.

Related