I have a 5 node cluster each having 132 GB physical memory and managed by Cloudera. However, both the properties yarn.nodemanager.resource.memory-mb and yarn.scheduler.maximum-allocation-mb have been set to 100 GB. I am submitting the job with 5 executors, 5 cores per executor and 40 GB of memory per executor. This code is just taking the file which is 16 GB in size and making it into RDD and caching it as follows
try {
BufferedReader br = new BufferedReader(new InputStreamReader(System.in));
System.out.println("Enter the file name");
String file = br.readLine();
// getting the context and spark configuration
SparkConf sparkConf = new SparkConf();
JavaSparkContext ctx = new JavaSparkContext(sparkConf);
// getting the input file into JavaRDD of strings
JavaRDD<String> lines = ctx.textFile("hdfs:/root/utkarsh/stark/" + file).cache();
lines.count();
ctx.close();
} catch (IOException e) {
e.printStackTrace();
}
While running I am observing the physical memory consumption of the 3 nodes running 2+2+1 executors respectively. The peak memory for the nodes are 54, 26, 48 GB respectively.
Why does the memory bloat whereas my input is only 16 GB.
Thanks in Advance!!