Class RemoteParForSparkWorker

java.lang.Object
org.apache.sysds.runtime.controlprogram.parfor.ParWorker
org.apache.sysds.runtime.controlprogram.parfor.RemoteParForSparkWorker
All Implemented Interfaces:
Serializable, org.apache.spark.api.java.function.PairFlatMapFunction<Task,Long,String>

public class RemoteParForSparkWorker extends ParWorker implements org.apache.spark.api.java.function.PairFlatMapFunction<Task,Long,String>
See Also:
  • Constructor Details

    • RemoteParForSparkWorker

      public RemoteParForSparkWorker(long jobid, String program, boolean isLocal, HashMap<String,byte[]> clsMap, boolean cpCaching, org.apache.spark.util.LongAccumulator atasks, org.apache.spark.util.LongAccumulator aiters, Map<String,org.apache.spark.broadcast.Broadcast<CacheBlock<?>>> brInputs, boolean cleanCache, Map<String,String> lineage)
  • Method Details

    • call

      public Iterator<scala.Tuple2<Long,String>> call(Task arg0) throws Exception
      Specified by:
      call in interface org.apache.spark.api.java.function.PairFlatMapFunction<Task,Long,String>
      Throws:
      Exception
    • cleanupCachedVariables

      public static void cleanupCachedVariables(long pfid)