Spark Java PCA: Java Heap Space e posizione di output mancante per shuffle

Oct 27 2020

Provo a fare un PCA su un dataframe con 4.827 righe e 40.107 colonne ma prendo un errore di spazio heap Java e posizione di output mancante per shuffle (secondo il file sdterr sugli eseguitori). L'errore si verifica durante la fase "treeAggregate su RowMatrix.scala: 122" della PCA.

Il cluster

Si tratta di un cluster autonomo con 16 nodi di lavoro, ciascuno con 1 esecutore con 4 core e 21,504 MB di memoria. Il nodo master ha 15g di memoria che io do con "Java -jar -Xmx15g myapp.jar". Anche "spark.sql.shuffle.partitions" sono 192 e "spark.driver.maxResultSize" è 6g.

Codice semplificato

df1.persist (From the Storage Tab in spark UI it says it is 3Gb)
df2=df1.groupby(col1).pivot(col2).mean(col3) (This is a df with 4.827 columns and 40.107 rows)
df2.collectFirstColumnAsList
df3=df1.groupby(col2).pivot(col1).mean(col3) (This is a df with 40.107 columns and 4.827 rows)

-----it hangs here for around 1.5 hours creating metadata for upcoming dataframe-----

df4 = (..Imputer or na.fill on df3..)
df5 = (..VectorAssembler on df4..)
(..PCA on df5 with error Missing output location for shuffle..)
df1.unpersist

Ho visto e provato molte soluzioni ma senza alcun risultato. Tra loro:

  1. Ripartizione di df5 o df4 a 16, 64, 192, 256, 1000, 4000 (anche se i dati non sembrano distorti)
  2. Modifica di spark.sql.shuffle.partitions su 16, 64, 192, 256, 1000, 4000
  3. Utilizzo di 1 e 2 core per esecutore in modo da avere più memoria per ogni attività.
  4. Avere 2 esecutori con 2 core o 4 core.
  5. Modifica "spark.memory.fraction" in 0.8 e "spark.memory.storageFraction" in 0.4.

Sempre lo stesso errore! Come è possibile spazzare via tutta questa memoria ?? È possibile che il df non si adatti alla memoria? Per favore fatemi sapere se avete bisogno di altre informazioni o schermate di stampa.

MODIFICA 1

Ho cambiato il cluster in 2 spark worker con 1 esecutore ciascuno con spark.sql.shuffle.partitions = 48. Ogni esecutore ha 115 ge 8 core. Di seguito è riportato il codice in cui carico il file (2.2Gb), converto ogni riga in un vettore denso e alimento il PCA.

Ogni riga del file ha questo formato (4.568 righe con 40.107 valori doppi ciascuna):

 "[x1,x2,x3,...]"

e il codice:

Dataset<Row> df1 = sp.read().format("com.databricks.spark.csv").option("header", "true").load("/home/ubuntu/yolo.csv");
StructType schema2 = new StructType(new StructField[] {
                        new StructField("intensity",new VectorUDT(),false,Metadata.empty())
            });
Dataset<Row> df = df1.map((Row originalrow) -> {
                    String yoho =originalrow.get(0).toString();
                    int sizeyoho=yoho.length();
                    String yohi = yoho.substring(1, sizeyoho-1);
                    String[] yi = yohi.split(",");
                    int s = yi.length;
                    double[] tmplist= new double[s];
                    for(int i=0;i<s;i++){
                        tmplist[i]=Double.parseDouble(yi[i]);
                    }
                    
                    Row newrow = RowFactory.create(Vectors.dense(tmplist));
                    return newrow;
            }, RowEncoder.apply(schema2));
PCAModel pcaexp = new PCA()
                    .setInputCol("intensity")
                    .setOutputCol("pcaFeatures")
                    .setK(2)
                    .fit(df);

L'errore esatto che ottengo sullo stderr di uno dei 2 worker è:

ERROR Executor: Exception in task 1.0 in stage 6.0 (TID 43)
java.lang.OutOfMemoryError
at java.io.ByteArrayOutputStream.hugeCapacity(ByteArrayOutputStream.java:123)
at java.io.ByteArrayOutputStream.grow(ByteArrayOutputStream.java:117)
at java.io.ByteArrayOutputStream.ensureCapacity(ByteArrayOutputStream.java:93)
at java.io.ByteArrayOutputStream.write(ByteArrayOutputStream.java:153)
at org.apache.spark.util.ByteBufferOutputStream.write(ByteBufferOutputStream.scala:41)
at java.io.ObjectOutputStream$BlockDataOutputStream.drain(ObjectOutputStream.java:1877) at java.io.ObjectOutputStream$BlockDataOutputStream.setBlockDataMode(ObjectOutputStream.java:1786)
at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1189)
at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348)
at org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:43)
at org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:100)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:456) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)

E questa è la scheda Stages di SparkUI:

E questo è lo stadio che fallisce (TreeAggregate su RowMatrix.scala: 122):

MODIFICA 2

MODIFICA 3

Ho letto l'intero file ma prendendo solo 10 valori da ogni riga e creando il vettore denso. Ricevo ancora lo stesso errore! Ho un master con 235 g di RAM e 3 lavoratori (1 esecutore ciascuno con 4 core) e 64 g di RAM per esecutore. Come è potuto succedere? (Non dimenticare che la dimensione totale del file è di soli 2,3 GB!)

Dataset<Row> df1 = sp.read().format("com.databricks.spark.csv").option("header", "true").load("/home/ubuntu/yolo.csv");

StructType schema2 = new StructType(new StructField[] {
                        new StructField("intensity",new VectorUDT(),false,Metadata.empty())
            });
Dataset<Row> df = df1.map((Row originalrow) -> {
                    String yoho =originalrow.get(0).toString();
                    int sizeyoho=yoho.length();
                    String yohi = yoho.substring(1, sizeyoho-1);
                    String[] yi = yohi.split(",");//this string array has all 40.107 values
                    int s = yi.length;
                    double[] tmplist= new double[s];
                    for(int i=0;i<10;i++){//I narrow it down to take only the first 10 values of each row
                        tmplist[i]=Double.parseDouble(yi[i]);
                    }
                    Row newrow = RowFactory.create(Vectors.dense(tmplist));
                    return newrow;
            }, RowEncoder.apply(schema2));
      
PCAModel pcaexp = new PCA()
                    .setInputCol("intensity")
                    .setOutputCol("pcaFeatures")
                    .setK(2)
                    .fit(df);

Risposte

1 EmiCareOfCell44 Oct 28 2020 at 10:14

La "posizione di output mancante per shuffle" si verifica quando la tua applicazione Spark esegue grandi fasi di shuffle, cerca di riallocare enormi quantità di dati tra gli esecutori e ci sono alcuni problemi nella tua rete cluster.

Spark dice che non hai memoria in qualche fase. Stai facendo trasformazioni che richiedono fasi diverse e consumano anche memoria. Inoltre, devi prima mantenere il dataframe e controllare il livello di archiviazione, perché è possibile che tu stia persistendo in memoria.

Stai concatenando diverse trasformazioni a livello di Spark: eseguendo il primo stadio pivot, ad esempio, Spark crea uno stadio ed esegue uno shuffle per raggruppare la tua colonna e forse hai dati disallineati e ci sono esecutori che consumano molta più memoria di altri, e forse l'errore può accadere in uno di essi.

Oltre alle trasformazioni Dataframe, lo stimatore PCA converte il dataframe in un RDD aumentando molto di più la memoria per calcolare la matrice covarianza, e lavora con rappresentazioni dense di matrici Breeze di NxN elementi non distribuiti . Ad esempio, l'SVD è realizzato con Breeze. Questo ha messo molta pressione su uno degli esecutori.

Forse puoi salvare il dataframe risultante in HDFS (o qualsiasi altra cosa) e fare la PCA un'altra applicazione Spark.

Il problema principale. quello che hai è che prima di de SVD l'algoritmo deve calcolare la matrice grammaticale e usa un treeAggregate da RDD. Questo crea una matrice Double molto grande che verrà inviata al driver e si verifica un errore perché il tuo driver non ha abbastanza memoria. È necessario aumentare notevolmente la memoria del driver. Hai errori di rete, se un esecutore perde la connessione il lavoro va in crash non tenta di rieseguire.

Personalmente, proverei a fare il PCA direttamente in Breeze (o Smile) nel driver, voglio dire, raccogliere il campo RDD perché il set di dati è abbastanza più piccolo della matrice covariante e farlo manualmente con una rappresentazione Float.

Codice per calcolare il PCA solo con Breeze, né Spark né TreeAgregation:

import breeze.linalg._
import breeze.linalg.svd._

object PCACode {
  
  def mean(v: Vector[Double]): Double = v.valuesIterator.sum / v.size

  def zeroMean(m: DenseMatrix[Double]): DenseMatrix[Double] = {
    val copy = m.copy
    for (c <- 0 until m.cols) {
      val col = copy(::, c)
      val colMean = mean(col)
      col -= colMean
    }
    copy
  }

  def pca(data: DenseMatrix[Double], components: Int): DenseMatrix[Double] = {
    val d = zeroMean(data)
    val SVD(_, _, v) = svd(d.t)
    val model = v(0 until components, ::)
    val filter = model.t * model
    filter * d
  }
  
  def main(args: Array[String]) : Unit = {
    val df : DataFrame = ???

    /** Collect the data and do the processing. Convert string to double, etc **/
    val data: Array[mutable.WrappedArray[Double]] =
      df.rdd.map(row => (row.getAs[mutable.WrappedArray[Double]](0))).collect()

    /** Once you have the Array, create the matrix and do the PCA **/
    val matrix = DenseMatrix(data.toSeq:_*)
    val pcaRes = pca(matrix, 2)

    println("result pca \n" + pcaRes)
  }
}

Questo codice eseguirà il PCA nel driver, controllerà la memoria. Se si blocca, potrebbe essere fatto con una precisione Float.