How persist(StorageLevel.MEMORY_AND_DISK) works in Spark 3.1 with Java implementation

Viewed 79

I am using Apache Spark 3.1 with java in GCP Dataproc cluster. And my code structure is like this.

Dataset<Row> dataset1 = readSpannerData(SparkSession session, Configuration session.sessionState().newHadoopConf());

Dataset<Row> dataset2 = reading some data from table1 bigtable

Dataset<Row> result1 = dataset1.join(dataset2);
dataset1.persist(StorageLevel.MEMORY_AND_DISK());
dataset2.persist(StorageLevel.MEMORY_AND_DISK()); //once the usage is done I am persisting both datasets

System.out.println(result1.count()); // It throws error in this line 

The exact error from YARN UI is, select query on spanner table which I am using in the starting of the job, not from any bigtable. I persisted dataset1 only after the usage is done.

And my cluster size is autoscale enabled with max of 250 worker nodes each have 8 core and 1024GB memory. It is configured to use 2 Executors on each node.(4 cores on each exe).

It was working fine with low volume of data. But it throws error while running with huge data.

Why it throws error in this situation? Will it look into the parent in memory dataset while using the result, calculated from the parent dataset which is already persisted? If we want to maintain that dataset then what is the usage of IN-Memory storage?

How it is working in low data environments? Howmany nodes and how long IN-MEMORY dataset will be maintained in spark job? Will the volume of the data affect IN-MEMORY dataset?

Can any one clarify this doubt?

Thanks In Advance :)

0 Answers
Related