Class DataPartitionerSparkMapper
java.lang.Object
org.apache.sysds.runtime.controlprogram.paramserv.dp.DataPartitionerSparkMapper
- All Implemented Interfaces:
Serializable,org.apache.spark.api.java.function.PairFlatMapFunction<scala.Tuple2<Long,scala.Tuple2<MatrixBlock, MatrixBlock>>, Integer, scala.Tuple2<Long, scala.Tuple2<MatrixBlock, MatrixBlock>>>
public class DataPartitionerSparkMapper
extends Object
implements org.apache.spark.api.java.function.PairFlatMapFunction<scala.Tuple2<Long,scala.Tuple2<MatrixBlock,MatrixBlock>>,Integer,scala.Tuple2<Long,scala.Tuple2<MatrixBlock,MatrixBlock>>>
- See Also:
-
Constructor Summary
ConstructorsConstructorDescriptionDataPartitionerSparkMapper(Statement.PSScheme scheme, int workersNum, SparkExecutionContext sec, int numEntries) -
Method Summary
Modifier and TypeMethodDescriptionIterator<scala.Tuple2<Integer,scala.Tuple2<Long, scala.Tuple2<MatrixBlock, MatrixBlock>>>> call(scala.Tuple2<Long, scala.Tuple2<MatrixBlock, MatrixBlock>> input) Do data partitioning
-
Constructor Details
-
DataPartitionerSparkMapper
public DataPartitionerSparkMapper(Statement.PSScheme scheme, int workersNum, SparkExecutionContext sec, int numEntries)
-
-
Method Details
-
call
public Iterator<scala.Tuple2<Integer,scala.Tuple2<Long, callscala.Tuple2<MatrixBlock, MatrixBlock>>>> (scala.Tuple2<Long, scala.Tuple2<MatrixBlock, throws ExceptionMatrixBlock>> input) Do data partitioning- Specified by:
callin interfaceorg.apache.spark.api.java.function.PairFlatMapFunction<scala.Tuple2<Long,scala.Tuple2<MatrixBlock, MatrixBlock>>, Integer, scala.Tuple2<Long, scala.Tuple2<MatrixBlock, MatrixBlock>>> - Parameters:
input- RowBlockID => (features, labels)- Returns:
- WorkerID => (rowBlockID, (single row features, single row labels))
- Throws:
Exception- Some exception
-