Class SparkUtils
java.lang.Object
org.apache.sysds.runtime.instructions.spark.utils.SparkUtils
-
Field Summary
Fields -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionstatic org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixCell> cacheBinaryCellRDD(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixCell> input) static voidcheckSparsity(String varname, ExecutionContext ec) static DataCharacteristicscomputeDataCharacteristics(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixCell> input) Utility to compute dimensions and non-zeros in a given RDD of binary cells.static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> copyBinaryBlockMatrix(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> in) Creates a partitioning-preserving deep copy of the input matrix RDD, where the indexes and values are copied.static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> copyBinaryBlockMatrix(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> in, boolean deep) Creates a partitioning-preserving copy of the input matrix RDD.static org.apache.spark.api.java.JavaPairRDD<TensorIndexes,BasicTensorBlock> copyBinaryBlockTensor(org.apache.spark.api.java.JavaPairRDD<TensorIndexes, BasicTensorBlock> in) Creates a partitioning-preserving deep copy of the input tensor RDD, where the indexes and values are copied.static org.apache.spark.api.java.JavaPairRDD<TensorIndexes,BasicTensorBlock> copyBinaryBlockTensor(org.apache.spark.api.java.JavaPairRDD<TensorIndexes, BasicTensorBlock> in, boolean deep) Creates a partitioning-preserving copy of the input tensor RDD.static List<scala.Tuple2<Long,FrameBlock>> static scala.Tuple2<Long,FrameBlock> static List<scala.Tuple2<MatrixIndexes,MatrixBlock>> static scala.Tuple2<MatrixIndexes,MatrixBlock> static List<Pair<MatrixIndexes,MatrixBlock>> static Pair<MatrixIndexes,MatrixBlock> static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> getEmptyBlockRDD(org.apache.spark.api.java.JavaSparkContext sc, DataCharacteristics mc) Creates an RDD of empty blocks according to the given matrix characteristics.static longgetNonZeros(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> input) static longstatic intstatic intgetNumPreferredPartitions(DataCharacteristics dc, boolean outputEmptyBlocks) static intgetNumPreferredPartitions(DataCharacteristics dc, org.apache.spark.api.java.JavaPairRDD<?, ?> in) static Stringstatic Stringstatic booleanisHashPartitioned(org.apache.spark.api.java.JavaPairRDD<?, ?> in) Indicates if the input RDD is hash partitioned, i.e., it has a partitioner of typeorg.apache.spark.HashPartitioner.static voidstatic Pair<Long,FrameBlock> toIndexedFrameBlock(scala.Tuple2<Long, FrameBlock> in) toIndexedLong(List<scala.Tuple2<Long, Long>> in) static IndexedMatrixValuestatic IndexedMatrixValuetoIndexedMatrixBlock(scala.Tuple2<MatrixIndexes, MatrixBlock> in) static IndexedTensorBlockstatic IndexedTensorBlocktoIndexedTensorBlock(scala.Tuple2<TensorIndexes, TensorBlock> in)
-
Field Details
-
DEFAULT_TMP
public static final org.apache.spark.storage.StorageLevel DEFAULT_TMP
-
-
Constructor Details
-
SparkUtils
public SparkUtils()
-
-
Method Details
-
toIndexedMatrixBlock
-
toIndexedMatrixBlock
-
toIndexedTensorBlock
-
toIndexedTensorBlock
-
fromIndexedMatrixBlock
-
fromIndexedMatrixBlock
public static List<scala.Tuple2<MatrixIndexes,MatrixBlock>> fromIndexedMatrixBlock(List<IndexedMatrixValue> in) -
fromIndexedMatrixBlockToPair
-
fromIndexedMatrixBlockToPair
public static List<Pair<MatrixIndexes,MatrixBlock>> fromIndexedMatrixBlockToPair(List<IndexedMatrixValue> in) -
fromIndexedFrameBlock
-
fromIndexedFrameBlock
public static List<scala.Tuple2<Long,FrameBlock>> fromIndexedFrameBlock(List<Pair<Long, FrameBlock>> in) -
toIndexedLong
-
toIndexedFrameBlock
-
isHashPartitioned
public static boolean isHashPartitioned(org.apache.spark.api.java.JavaPairRDD<?, ?> in) Indicates if the input RDD is hash partitioned, i.e., it has a partitioner of typeorg.apache.spark.HashPartitioner.- Parameters:
in- input JavaPairRDD- Returns:
- true if input is hash partitioned
-
getNumPreferredPartitions
public static int getNumPreferredPartitions(DataCharacteristics dc, org.apache.spark.api.java.JavaPairRDD<?, ?> in) -
getNumPreferredPartitions
-
getNumPreferredPartitions
-
copyBinaryBlockMatrix
public static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> copyBinaryBlockMatrix(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> in) Creates a partitioning-preserving deep copy of the input matrix RDD, where the indexes and values are copied.- Parameters:
in- matrix asJavaPairRDD<MatrixIndexes,MatrixBlock>- Returns:
- matrix as
JavaPairRDD<MatrixIndexes,MatrixBlock>
-
copyBinaryBlockMatrix
public static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> copyBinaryBlockMatrix(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> in, boolean deep) Creates a partitioning-preserving copy of the input matrix RDD. If a deep copy is requested, indexes and values are copied, otherwise they are simply passed through.- Parameters:
in- matrix asJavaPairRDD<MatrixIndexes,MatrixBlock>deep- if true, perform deep copy- Returns:
- matrix as
JavaPairRDD<MatrixIndexes,MatrixBlock>
-
copyBinaryBlockTensor
public static org.apache.spark.api.java.JavaPairRDD<TensorIndexes,BasicTensorBlock> copyBinaryBlockTensor(org.apache.spark.api.java.JavaPairRDD<TensorIndexes, BasicTensorBlock> in) Creates a partitioning-preserving deep copy of the input tensor RDD, where the indexes and values are copied.- Parameters:
in- tensor asJavaPairRDD<TensorIndexes,HomogTensor>- Returns:
- tensor as
JavaPairRDD<TensorIndexes,HomogTensor>
-
copyBinaryBlockTensor
public static org.apache.spark.api.java.JavaPairRDD<TensorIndexes,BasicTensorBlock> copyBinaryBlockTensor(org.apache.spark.api.java.JavaPairRDD<TensorIndexes, BasicTensorBlock> in, boolean deep) Creates a partitioning-preserving copy of the input tensor RDD. If a deep copy is requested, indexes and values are copied, otherwise they are simply passed through.- Parameters:
in- tensor asJavaPairRDD<TensorIndexes,HomogTensor>deep- if true, perform deep copy- Returns:
- tensor as
JavaPairRDD<TensorIndexes,HomogTensor>
-
checkSparsity
-
getStartLineFromSparkDebugInfo
-
getPrefixFromSparkDebugInfo
-
getEmptyBlockRDD
public static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> getEmptyBlockRDD(org.apache.spark.api.java.JavaSparkContext sc, DataCharacteristics mc) Creates an RDD of empty blocks according to the given matrix characteristics. This is done in a scalable manner by parallelizing block ranges and generating empty blocks in a distributed manner, under awareness of preferred output partition sizes.- Parameters:
sc- spark contextmc- matrix characteristics- Returns:
- pair rdd of empty matrix blocks
-
cacheBinaryCellRDD
public static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixCell> cacheBinaryCellRDD(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixCell> input) -
computeDataCharacteristics
public static DataCharacteristics computeDataCharacteristics(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixCell> input) Utility to compute dimensions and non-zeros in a given RDD of binary cells.- Parameters:
input- matrix asJavaPairRDD<MatrixIndexes, MatrixCell>- Returns:
- matrix characteristics
-
getNonZeros
-
getNonZeros
public static long getNonZeros(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> input) -
postprocessUltraSparseOutput
-