Переглянути джерело

1.Add num-executor option in BatchRequest; 2. Fix bug of files option

Tianjin Gu 10 роки тому
батько
коміт
74d6181188

+ 12 - 1
apps/spark/java/livy-core/src/main/scala/com/cloudera/hue/livy/spark/SparkSubmitProcessBuilder.scala

@@ -181,6 +181,16 @@ class SparkSubmitProcessBuilder(livyConf: LivyConf) extends Logging {
     this
   }
 
+
+  def numExecutors(numExecutors: Int): SparkSubmitProcessBuilder = {
+    this.numExecutors(numExecutors.toString)
+  }
+
+  def numExecutors(numExecutors: String): SparkSubmitProcessBuilder = {
+    _numExecutors = Some(numExecutors)
+    this
+  }
+
   def queue(queue: String): SparkSubmitProcessBuilder = {
     _queue = Some(queue)
     this
@@ -249,6 +259,7 @@ class SparkSubmitProcessBuilder(livyConf: LivyConf) extends Logging {
     addOpt("--proxy-user", _proxyUser)
     addOpt("--driver-cores", _driverCores)
     addOpt("--executor-cores", _executorCores)
+    addOpt("--num-executors", _numExecutors)
     addOpt("--queue", _queue)
     addList("--archives", _archives.map(fromPath))
 
@@ -274,7 +285,7 @@ class SparkSubmitProcessBuilder(livyConf: LivyConf) extends Logging {
   private def fromPath(path: Path) = path match {
     case AbsolutePath(p) => p
     case RelativePath(p) =>
-      if (p.startsWith("hdfs://")) {
+      if (p.startsWith("hdfs://") || p.startsWith("file://")) {
         p
       } else {
         fsRoot + "/" + p

+ 1 - 0
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/batch/BatchSessionManager.scala

@@ -65,4 +65,5 @@ case class CreateBatchRequest(file: String,
                               driverCores: Option[Int] = None,
                               executorMemory: Option[String] = None,
                               executorCores: Option[Int] = None,
+                              numExecutors: Option[Int] = None,
                               archives: List[String] = List())

+ 1 - 0
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/batch/BatchSessionYarn.scala

@@ -59,6 +59,7 @@ object BatchSessionYarn {
     createBatchRequest.driverCores.foreach(builder.driverCores)
     createBatchRequest.executorMemory.foreach(builder.executorMemory)
     createBatchRequest.executorCores.foreach(builder.executorCores)
+    createBatchRequest.numExecutors.foreach(builder.numExecutors)
     createBatchRequest.archives.map(RelativePath).foreach(builder.archive)
 
     builder.redirectOutput(Redirect.PIPE)