Spark Java PCA: Java Heap Space và Thiếu vị trí đầu ra để trộn
Tôi cố gắng thực hiện PCA trên khung dữ liệu có 4,827 hàng và 40,107 cột nhưng tôi gặp lỗi không gian đống Java và thiếu vị trí đầu ra để trộn (theo tệp sdterr trên trình thực thi). Lỗi xảy ra trong giai đoạn "treeAggregate at RowMatrix.scala: 122" của PCA.
Cụm
Nó là một cụm độc lập với 16 nút công nhân, mỗi nút có 1 người thực thi với 4 lõi và bộ nhớ 21.504mb. Nút chính có bộ nhớ 15g mà tôi cung cấp với "Java -jar -Xmx15g myapp.jar". Ngoài ra "spark.sql.shuffle.partitions" là 192 và "spark.driver.maxResultSize" là 6g.
Mã đơn giản
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
Tôi đã xem và thử nhiều giải pháp nhưng không có kết quả. Trong số đó:
- Phân vùng lại df5 hoặc df4 thành 16, 64, 192, 256, 1000, 4000 (mặc dù dữ liệu không bị lệch)
- Thay đổi tiêu đề spark.sql.shuffle.partitions thành 16, 64, 192, 256, 1000, 4000
- Sử dụng 1 và 2 lõi cho mỗi trình thực thi để có nhiều bộ nhớ hơn cho mọi tác vụ.
- Có 2 người thực thi với 2 lõi hoặc 4 lõi.
- Thay đổi "spark.memory.fraction" thành 0,8 và "spark.memory.storageFraction" thành 0,4.
Luôn luôn cùng một lỗi! Làm thế nào để có thể thổi bay tất cả ký ức này ?? Có thể df thực sự không phù hợp trong bộ nhớ? Vui lòng cho tôi biết nếu bạn cần bất kỳ thông tin hoặc màn hình in nào khác.
CHỈNH SỬA 1
Tôi đã thay đổi cụm thành 2 công nhân tia lửa với 1 người thực thi mỗi người với spark.sql.shuffle.partitions = 48. Mỗi viên thực thi có 115g và 8 lõi. Dưới đây là mã nơi tôi tải tệp (2.2Gb), chuyển đổi mỗi dòng thành một vectơ dày đặc và cung cấp cho PCA.
Mỗi hàng trong tệp có định dạng này (4,568 hàng với 40,107 giá trị kép mỗi hàng):
"[x1,x2,x3,...]"
và mã:
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ỗi chính xác tôi nhận được trên stderr của một trong 2 công nhân là:
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)
Và đây là Tab các giai đoạn của SparkUI:
Và đây là Giai đoạn bị lỗi (TreeAggregate tại RowMatrix.scala: 122):
CHỈNH SỬA 2
CHỈNH SỬA 3
Tôi đọc toàn bộ tệp nhưng chỉ lấy 10 giá trị từ mỗi hàng và tạo vectơ dày đặc. Tôi vẫn gặp lỗi tương tự! Tôi có một bậc thầy với 235g Ram và 3 công nhân (1 người thi hành mỗi người có 4 lõi) và 64g Ram cho mỗi người thực thi. Làm thế nào điều này có thể xảy ra? (Đừng quên tổng kích thước của tệp chỉ là 2.3Gb!)
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);
Trả lời
Các "Thiếu vị trí đầu ra cho shuffle" xảy ra khi ứng dụng Spark bạn làm giai đoạn xáo trộn lớn, nó sẽ cố gắng để tái phân bổ số lượng lớn các dữ liệu giữa các chấp hành viên và có một số vấn đề trong mạng cluster của bạn.
Spark nói rằng bạn không có trí nhớ trong một số giai đoạn. Bạn đang thực hiện các phép biến đổi đòi hỏi các giai đoạn khác nhau và chúng cũng tiêu tốn bộ nhớ. Bên cạnh đó, trước tiên bạn vẫn duy trì khung dữ liệu và nên kiểm tra mức lưu trữ, vì có thể bạn đang duy trì trong bộ nhớ.
Bạn đang xâu chuỗi một số phép biến đổi rộng Spark: thực hiện giai đoạn xoay vòng đầu tiên, ví dụ: Spark tạo một giai đoạn và thực hiện xáo trộn để nhóm cho cột của bạn và có thể bạn bị lệch dữ liệu và có những trình thực thi tiêu tốn nhiều bộ nhớ hơn những người khác, và có thể lỗi có thể xảy ra ở một trong số chúng.
Bên cạnh các phép biến đổi Khung dữ liệu, công cụ ước tính PCA chuyển đổi khung dữ liệu thành RDD làm tăng nhiều bộ nhớ hơn để tính toán ma trận covarianze và nó hoạt động với các đại diện dày đặc của ma trận Breeze của các phần tử NxN không được phân phối . Ví dụ, SVD được tạo bằng Breeze. Điều đó gây ra rất nhiều áp lực cho một trong những người thực thi.
Có thể bạn có thể lưu khung dữ liệu kết quả trong HDFS (hoặc bất cứ thứ gì) và thực hiện PCA một ứng dụng Spark khác.
Vấn đề chính. mà bạn có là trước khi de SVD, thuật toán cần tính toán Ma trận Grammian và nó sử dụng một treeAggregate từ RDD. Điều này tạo ra một ma trận kép rất lớn sẽ được gửi đến trình điều khiển và có lỗi do trình điều khiển của bạn không có đủ bộ nhớ. Bạn cần tăng đáng kể bộ nhớ trình điều khiển. Bạn có lỗi mạng, nếu một người thực thi mất kết nối, công việc bị treo, nó sẽ không thử thực thi lại.
Cá nhân, tôi sẽ cố gắng thực hiện PCA trực tiếp trong Breeze (hoặc Smile) trong trình điều khiển, ý tôi là, thu thập trường RDD vì tập dữ liệu khá nhỏ hơn ma trận covarianze và thực hiện thủ công với biểu diễn Float.
Mã để tính toán PCA chỉ với Breeze, không phải Spark hay 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)
}
}
Mã này sẽ thực hiện PCA trong trình điều khiển, kiểm tra bộ nhớ. Nếu nó gặp sự cố, nó có thể được thực hiện với một precission Float.