Class SparkExecutionContext
java.lang.Object
org.apache.sysds.runtime.controlprogram.context.ExecutionContext
org.apache.sysds.runtime.controlprogram.context.SparkExecutionContext
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic classstatic classCaptures relevant spark cluster configuration properties, e.g., memory budgets and degree of parallelism. -
Field Summary
Fields -
Method Summary
Modifier and TypeMethodDescriptionvoidaddLineage(String varParent, String varChild, boolean broadcast) voidaddLineageBroadcast(String varParent, String varChild) Adds a child broadcast object to the lineage of a parent rdd.voidaddLineageRDD(String varParent, String varChild) Adds a child rdd object to the lineage of a parent rdd.org.apache.spark.broadcast.Broadcast<CacheBlock<?>>broadcastVariable(CacheableData<CacheBlock<?>> cd) voidcacheMatrixObject(String var) static voidcleanupBroadcastVariable(org.apache.spark.broadcast.Broadcast<?> bvar) This call destroys a broadcast variable at all executors and the driver.voidstatic voidcleanupRDDVariable(org.apache.spark.api.java.JavaPairRDD<?, ?> rvar) This call removes an rdd variable from executor memory and disk if required.static voidvoidcleanupThreadLocalSchedulerPool(int pool) voidclose()static org.apache.spark.SparkConfSets up a SystemDS-preferred Spark configuration based on the implicit default configuration (as passed via configurations from outside).static voidRegisters an active user of the shared spark context.static voidReleases an active user previously registered viaenterSparkExecution().org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> Spark instructions should call this for all matrix inputs except broadcast variables.org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> getBinaryMatrixBlockRDDHandleForVariable(String varname, int numParts, boolean inclEmpty) org.apache.spark.api.java.JavaPairRDD<TensorIndexes,TensorBlock> Spark instructions should call this for all tensor inputs except broadcast variables.org.apache.spark.api.java.JavaPairRDD<TensorIndexes,TensorBlock> getBinaryTensorBlockRDDHandleForVariable(String varname, int numParts, boolean inclEmpty) getBroadcastForFrameVariable(String varname) getBroadcastForTensorVariable(String varname) getBroadcastForVariable(String varname) static doubleObtains the available memory budget for broadcast variables in bytes.static doublegetDataMemoryBudget(boolean min, boolean refresh) Obtain the available memory budget for data storage in bytes.static intgetDefaultParallelism(boolean refresh) Obtain the default degree of parallelism (cores in the cluster).org.apache.spark.api.java.JavaPairRDD<Long,FrameBlock> Spark instructions should call this for all frame inputs except broadcast variables.static longgetMemCachedRDDSize(int rddID) static intObtain the number of executors in the cluster (excluding the driver).org.apache.spark.api.java.JavaPairRDD<?,?> FIXME: currently this implementation assumes matrix representations but frame signature in order to support the old transform implementation.org.apache.spark.api.java.JavaPairRDD<?,?> org.apache.spark.api.java.JavaPairRDD<?,?> getRDDHandleForMatrixObject(MatrixObject mo, Types.FileFormat fmt, int numParts, boolean inclEmpty) org.apache.spark.api.java.JavaPairRDD<?,?> getRDDHandleForTensorObject(TensorObject to, Types.FileFormat fmt, int numParts, boolean inclEmpty) org.apache.spark.api.java.JavaPairRDD<?,?> getRDDHandleForVariable(String varname, Types.FileFormat fmt, int numParts, boolean inclEmpty) Obtains the lazily analyzed spark cluster configuration.org.apache.spark.api.java.JavaSparkContextReturns the used singleton spark context.static org.apache.spark.api.java.JavaSparkContextstatic longstatic voidstatic voidinitLocalSparkContext(org.apache.spark.SparkConf sparkConf) static voidstatic booleanstatic booleanstatic booleanisRDDCached(int rddID) static booleanIndicates if the spark context has been created or has been passed in from outside.voidstatic voidvoidvoidsetRDDHandleForVariable(String varname, org.apache.spark.api.java.JavaPairRDD<?, ?> rdd) Keep the output rdd of spark rdd operations as meta data of matrix/frame objects in the symbol table.voidsetRDDHandleForVariable(String varname, RDDObject rddhandle) intstatic FrameBlocktoFrameBlock(org.apache.spark.api.java.JavaPairRDD<Long, FrameBlock> rdd, Types.ValueType[] schema, int rlen, int clen) static FrameBlocktoFrameBlock(RDDObject rdd, Types.ValueType[] schema, int rlen, int clen) static org.apache.spark.api.java.JavaPairRDD<Long,FrameBlock> toFrameJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, FrameBlock src) static MatrixBlocktoMatrixBlock(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> rdd, int rlen, int clen, int blen, long nnz) Utility method for creating a single matrix block out of a binary block RDD.static MatrixBlocktoMatrixBlock(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixCell> rdd, int rlen, int clen, long nnz) Utility method for creating a single matrix block out of a binary cell RDD.static MatrixBlocktoMatrixBlock(RDDObject rdd, int rlen, int clen, int blen, long nnz) This method is a generic abstraction for calls from the buffer pool.static MatrixBlocktoMatrixBlock(RDDObject rdd, int rlen, int clen, long nnz) static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> toMatrixJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, MatrixBlock src, int blen) static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> toMatrixJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, MatrixBlock src, int blen, int numParts, boolean inclEmpty) static PartitionedBlock<MatrixBlock>toPartitionedMatrixBlock(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> rdd, int rlen, int clen, int blen, long nnz) static TensorBlocktoTensorBlock(org.apache.spark.api.java.JavaPairRDD<TensorIndexes, TensorBlock> rdd, DataCharacteristics dc) static org.apache.spark.api.java.JavaPairRDD<TensorIndexes,TensorBlock> toTensorJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, TensorBlock src, int blen) static org.apache.spark.api.java.JavaPairRDD<TensorIndexes,TensorBlock> toTensorJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, TensorBlock src, int blen, int numParts, boolean inclEmpty) static voidwriteFrameRDDtoHDFS(RDDObject rdd, String path, Types.FileFormat fmt) static longwriteMatrixRDDtoHDFS(RDDObject rdd, String path, Types.FileFormat fmt) Methods inherited from class org.apache.sysds.runtime.controlprogram.context.ExecutionContext
addTmpParforFunction, allocateGPUMatrixObject, cleanupDataObject, containsVariable, containsVariable, createCacheableData, createFrameObject, createFrameObject, createMatrixObject, createMatrixObject, getCacheableData, getCacheableData, getDataCharacteristics, getDenseMatrixOutputForGPUInstruction, getDenseMatrixOutputForGPUInstruction, getFrameInput, getFrameInput, getFrameInputs, getFrameObject, getFrameObject, getGPUContext, getGPUContexts, getGPUDensePointerAddress, getGPUSparsePointerAddress, getLineage, getLineageItem, getLineageItem, getListObject, getListObject, getMatrixInput, getMatrixInput, getMatrixInputForGPUInstruction, getMatrixInputs, getMatrixInputs, getMatrixLineagePair, getMatrixLineagePair, getMatrixObject, getMatrixObject, getMetaData, getNumGPUContexts, getOrCreateLineageItem, getProgram, getScalarInput, getScalarInput, getScalarInputs, getSealClient, getSparseMatrixOutputForGPUInstruction, getSparseMatrixOutputForGPUInstruction, getTensorInput, getTensorObject, getTID, getTmpParforFunctions, getVariable, getVariable, getVariables, getVarList, getVarListPartitioned, isAutoCreateVars, isFederated, isFederated, isFrameObject, isMatrixObject, maintainLineageDebuggerInfo, pinVariables, releaseCacheableData, releaseFrameInput, releaseFrameInputs, releaseMatrixInput, releaseMatrixInput, releaseMatrixInputForGPUInstruction, releaseMatrixInputs, releaseMatrixInputs, releaseMatrixOutputForGPUInstruction, releaseTensorInput, releaseTensorInput, removeVariable, replaceLineageItem, setAutoCreateVars, setFrameOutput, setGPUContexts, setLineage, setMatrixOutput, setMatrixOutput, setMatrixOutput, setMatrixOutputAndLineage, setMatrixOutputAndLineage, setMatrixOutputAndLineage, setMetaData, setMetaData, setProgram, setScalarOutput, setSealClient, setTensorOutput, setTID, setVariable, setVariables, toString, traceLineage, unpinVariables
-
Field Details
-
FAIR_SCHEDULER_MODE
public static final boolean FAIR_SCHEDULER_MODE- See Also:
-
-
Method Details
-
getSparkContext
public org.apache.spark.api.java.JavaSparkContext getSparkContext()Returns the used singleton spark context. In case of lazy spark context creation, this methods blocks until the spark context is created.- Returns:
- java spark context
-
initLocalSparkContext
public static void initLocalSparkContext(org.apache.spark.SparkConf sparkConf) -
getSparkContextStatic
public static org.apache.spark.api.java.JavaSparkContext getSparkContextStatic() -
isSparkContextCreated
public static boolean isSparkContextCreated()Indicates if the spark context has been created or has been passed in from outside.- Returns:
- true if spark context created
-
resetSparkContextStatic
public static void resetSparkContextStatic() -
enterSparkExecution
public static void enterSparkExecution()Registers an active user of the shared spark context. Must be balanced by a laterexitSparkExecution()so a concurrent execution cannot stop the context while this one still has in-flight jobs. -
exitSparkExecution
public static void exitSparkExecution()Releases an active user previously registered viaenterSparkExecution(). Only adjusts the count; the actual teardown is left toclose(), which stops the context once no registered execution remains. -
close
public void close() -
isLazySparkContextCreation
public static boolean isLazySparkContextCreation() -
handleIllegalReflectiveAccessSpark
public static void handleIllegalReflectiveAccessSpark() -
initSparkContext
public static void initSparkContext() -
createSystemDSSparkConf
public static org.apache.spark.SparkConf createSystemDSSparkConf()Sets up a SystemDS-preferred Spark configuration based on the implicit default configuration (as passed via configurations from outside).- Returns:
- spark configuration
-
isLocalMaster
public static boolean isLocalMaster() -
getBinaryMatrixBlockRDDHandleForVariable
public org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> getBinaryMatrixBlockRDDHandleForVariable(String varname) Spark instructions should call this for all matrix inputs except broadcast variables.- Parameters:
varname- variable name- Returns:
- JavaPairRDD of MatrixIndexes-MatrixBlocks
-
getBinaryMatrixBlockRDDHandleForVariable
public org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> getBinaryMatrixBlockRDDHandleForVariable(String varname, int numParts, boolean inclEmpty) -
getBinaryTensorBlockRDDHandleForVariable
public org.apache.spark.api.java.JavaPairRDD<TensorIndexes,TensorBlock> getBinaryTensorBlockRDDHandleForVariable(String varname) Spark instructions should call this for all tensor inputs except broadcast variables.- Parameters:
varname- variable name- Returns:
- JavaPairRDD of TensorIndexes-HomogTensors
-
getBinaryTensorBlockRDDHandleForVariable
public org.apache.spark.api.java.JavaPairRDD<TensorIndexes,TensorBlock> getBinaryTensorBlockRDDHandleForVariable(String varname, int numParts, boolean inclEmpty) -
getFrameBinaryBlockRDDHandleForVariable
public org.apache.spark.api.java.JavaPairRDD<Long,FrameBlock> getFrameBinaryBlockRDDHandleForVariable(String varname) Spark instructions should call this for all frame inputs except broadcast variables.- Parameters:
varname- variable name- Returns:
- JavaPairRDD of Longs-FrameBlocks
-
getRDDHandleForVariable
public org.apache.spark.api.java.JavaPairRDD<?,?> getRDDHandleForVariable(String varname, Types.FileFormat fmt, int numParts, boolean inclEmpty) -
getRDDHandleForMatrixObject
public org.apache.spark.api.java.JavaPairRDD<?,?> getRDDHandleForMatrixObject(MatrixObject mo, Types.FileFormat fmt) -
getRDDHandleForMatrixObject
public org.apache.spark.api.java.JavaPairRDD<?,?> getRDDHandleForMatrixObject(MatrixObject mo, Types.FileFormat fmt, int numParts, boolean inclEmpty) -
getRDDHandleForTensorObject
public org.apache.spark.api.java.JavaPairRDD<?,?> getRDDHandleForTensorObject(TensorObject to, Types.FileFormat fmt, int numParts, boolean inclEmpty) -
getRDDHandleForFrameObject
public org.apache.spark.api.java.JavaPairRDD<?,?> getRDDHandleForFrameObject(FrameObject fo, Types.FileFormat fmt) FIXME: currently this implementation assumes matrix representations but frame signature in order to support the old transform implementation.- Parameters:
fo- frame objectfmt- file format type- Returns:
- JavaPairRDD handle for a frame object
-
broadcastVariable
public org.apache.spark.broadcast.Broadcast<CacheBlock<?>> broadcastVariable(CacheableData<CacheBlock<?>> cd) -
getBroadcastForMatrixObject
-
setBroadcastHandle
-
getBroadcastForTensorObject
-
getBroadcastForVariable
-
getBroadcastForTensorVariable
-
getBroadcastForFrameVariable
-
setRDDHandleForVariable
Keep the output rdd of spark rdd operations as meta data of matrix/frame objects in the symbol table.- Parameters:
varname- variable namerdd- JavaPairRDD handle for variable
-
setRDDHandleForVariable
-
toMatrixJavaPairRDD
public static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> toMatrixJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, MatrixBlock src, int blen) -
toMatrixJavaPairRDD
public static org.apache.spark.api.java.JavaPairRDD<MatrixIndexes,MatrixBlock> toMatrixJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, MatrixBlock src, int blen, int numParts, boolean inclEmpty) -
toTensorJavaPairRDD
public static org.apache.spark.api.java.JavaPairRDD<TensorIndexes,TensorBlock> toTensorJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, TensorBlock src, int blen) -
toTensorJavaPairRDD
public static org.apache.spark.api.java.JavaPairRDD<TensorIndexes,TensorBlock> toTensorJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, TensorBlock src, int blen, int numParts, boolean inclEmpty) -
toFrameJavaPairRDD
public static org.apache.spark.api.java.JavaPairRDD<Long,FrameBlock> toFrameJavaPairRDD(org.apache.spark.api.java.JavaSparkContext sc, FrameBlock src) -
toMatrixBlock
This method is a generic abstraction for calls from the buffer pool.- Parameters:
rdd- rdd objectrlen- number of rowsclen- number of columnsblen- block lengthnnz- number of non-zeros- Returns:
- matrix block
-
toMatrixBlock
public static MatrixBlock toMatrixBlock(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> rdd, int rlen, int clen, int blen, long nnz) Utility method for creating a single matrix block out of a binary block RDD. Note that this collect call might trigger execution of any pending transformations. NOTE: This is an unguarded utility function, which requires memory for both the output matrix and its collected, blocked representation.- Parameters:
rdd- JavaPairRDD for matrix blockrlen- number of rowsclen- number of columnsblen- block lengthnnz- number of non-zeros- Returns:
- Local matrix block
-
toMatrixBlock
-
toMatrixBlock
public static MatrixBlock toMatrixBlock(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixCell> rdd, int rlen, int clen, long nnz) Utility method for creating a single matrix block out of a binary cell RDD. Note that this collect call might trigger execution of any pending transformations.- Parameters:
rdd- JavaPairRDD for matrix blockrlen- number of rowsclen- number of columnsnnz- number of non-zeros- Returns:
- matrix block
-
toTensorBlock
public static TensorBlock toTensorBlock(org.apache.spark.api.java.JavaPairRDD<TensorIndexes, TensorBlock> rdd, DataCharacteristics dc) -
toPartitionedMatrixBlock
public static PartitionedBlock<MatrixBlock> toPartitionedMatrixBlock(org.apache.spark.api.java.JavaPairRDD<MatrixIndexes, MatrixBlock> rdd, int rlen, int clen, int blen, long nnz) -
toFrameBlock
-
toFrameBlock
public static FrameBlock toFrameBlock(org.apache.spark.api.java.JavaPairRDD<Long, FrameBlock> rdd, Types.ValueType[] schema, int rlen, int clen) -
writeMatrixRDDtoHDFS
-
writeFrameRDDtoHDFS
-
addLineageRDD
Adds a child rdd object to the lineage of a parent rdd.- Parameters:
varParent- parent variablevarChild- child variable
-
addLineageBroadcast
Adds a child broadcast object to the lineage of a parent rdd.- Parameters:
varParent- parent variablevarChild- child variable
-
addLineage
-
cleanupCacheableData
- Overrides:
cleanupCacheableDatain classExecutionContext
-
cleanupSingleLineageObject
-
cleanupBroadcastVariable
public static void cleanupBroadcastVariable(org.apache.spark.broadcast.Broadcast<?> bvar) This call destroys a broadcast variable at all executors and the driver. Hence, it is intended to be used on rmvar only. Depending on the ASYNCHRONOUS_VAR_DESTROY configuration, this is asynchronous or not.- Parameters:
bvar- broadcast variable
-
cleanupRDDVariable
public static void cleanupRDDVariable(org.apache.spark.api.java.JavaPairRDD<?, ?> rvar) This call removes an rdd variable from executor memory and disk if required. Hence, it is intended to be used on rmvar only. Depending on the ASYNCHRONOUS_VAR_DESTROY configuration, this is asynchronous or not.- Parameters:
rvar- rdd variable to remove
-
repartitionAndCacheMatrixObject
-
cacheMatrixObject
-
setThreadLocalSchedulerPool
public int setThreadLocalSchedulerPool() -
cleanupThreadLocalSchedulerPool
public void cleanupThreadLocalSchedulerPool(int pool) -
isRDDCached
public static boolean isRDDCached(int rddID) -
getMemCachedRDDSize
public static long getMemCachedRDDSize(int rddID) -
getStorageSpaceUsed
public static long getStorageSpaceUsed() -
getSparkClusterConfig
Obtains the lazily analyzed spark cluster configuration.- Returns:
- spark cluster configuration
-
getBroadcastMemoryBudget
public static double getBroadcastMemoryBudget()Obtains the available memory budget for broadcast variables in bytes.- Returns:
- broadcast memory budget
-
getDataMemoryBudget
public static double getDataMemoryBudget(boolean min, boolean refresh) Obtain the available memory budget for data storage in bytes.- Parameters:
min- flag for minimum data budgetrefresh- flag for refresh with spark context- Returns:
- data memory budget
-
getNumExecutors
public static int getNumExecutors()Obtain the number of executors in the cluster (excluding the driver).- Returns:
- number of executors
-
getDefaultParallelism
public static int getDefaultParallelism(boolean refresh) Obtain the default degree of parallelism (cores in the cluster).- Parameters:
refresh- flag for refresh with spark context- Returns:
- default degree of parallelism
-