Преглед изворни кода

[livy] Allow interactive session to upload jars and files into the working dir

Erick Tryzelaar пре 10 година
родитељ
комит
6c0f1d7

+ 5 - 1
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/interactive/InteractiveSessionManager.scala

@@ -99,10 +99,14 @@ class SessionManager(factory: InteractiveSessionFactory) extends Logging {
 
 case class CreateInteractiveRequest(kind: Kind,
                                     proxyUser: Option[String] = None,
+                                    jars: List[String] = List(),
+                                    pyFiles: List[String] = List(),
+                                    files: List[String] = List(),
                                     driverMemory: Option[String] = None,
                                     driverCores: Option[Int] = None,
                                     executorMemory: Option[String] = None,
-                                    executorCores: Option[Int] = None)
+                                    executorCores: Option[Int] = None,
+                                    archives: List[String] = List())
 
 class SessionNotFound extends Exception
 

+ 7 - 4
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/interactive/InteractiveSessionProcess.scala

@@ -22,7 +22,7 @@ import java.lang.ProcessBuilder.Redirect
 import java.net.URL
 
 import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder
-import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.AbsolutePath
+import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.{RelativePath, AbsolutePath}
 import com.cloudera.hue.livy.{LivyConf, Logging, Utils}
 
 import scala.annotation.tailrec
@@ -45,16 +45,19 @@ object InteractiveSessionProcess extends Logging {
 
     val builder = new SparkSubmitProcessBuilder(livyConf)
 
+    builder.className("com.cloudera.hue.livy.repl.Main")
+    createInteractiveRequest.archives.map(RelativePath).foreach(builder.archive)
     createInteractiveRequest.driverCores.foreach(builder.driverCores)
     createInteractiveRequest.driverMemory.foreach(builder.driverMemory)
     createInteractiveRequest.executorCores.foreach(builder.executorCores)
     createInteractiveRequest.executorMemory.foreach(builder.executorMemory)
-
-    builder.className("com.cloudera.hue.livy.repl.Main")
+    createInteractiveRequest.files.map(RelativePath).foreach(builder.file)
+    createInteractiveRequest.jars.map(RelativePath).foreach(builder.jar)
+    createInteractiveRequest.proxyUser.foreach(builder.proxyUser)
+    createInteractiveRequest.pyFiles.map(RelativePath).foreach(builder.pyFile)
 
     sys.env.get("LIVY_REPL_JAVA_OPTS").foreach(builder.driverJavaOptions)
     livyConf.getOption(CONF_LIVY_REPL_DRIVER_CLASS_PATH).foreach(builder.driverClassPath)
-    createInteractiveRequest.proxyUser.foreach(builder.proxyUser)
 
     livyConf.getOption(CONF_LIVY_REPL_CALLBACK_URL).foreach { case callbackUrl =>
       builder.env("LIVY_CALLBACK_URL", f"$callbackUrl/sessions/$id/callback")

+ 8 - 4
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/interactive/InteractiveSessionYarn.scala

@@ -23,7 +23,7 @@ import java.util.concurrent.TimeUnit
 
 import com.cloudera.hue.livy.sessions.Error
 import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder
-import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.AbsolutePath
+import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.{RelativePath, AbsolutePath}
 import com.cloudera.hue.livy.yarn.{Client, Job}
 import com.cloudera.hue.livy.{LineBufferedProcess, LivyConf, Utils}
 
@@ -45,11 +45,15 @@ object InteractiveSessionYarn {
     builder.master("yarn-cluster")
     builder.className("com.cloudera.hue.livy.repl.Main")
     builder.driverJavaOptions(f"-Dlivy.repl.callback-url=$url -Dlivy.repl.port=0")
-    createInteractiveRequest.proxyUser.foreach(builder.proxyUser)
-    createInteractiveRequest.driverMemory.foreach(builder.driverMemory)
+    createInteractiveRequest.archives.map(RelativePath).foreach(builder.archive)
     createInteractiveRequest.driverCores.foreach(builder.driverCores)
-    createInteractiveRequest.executorMemory.foreach(builder.executorMemory)
+    createInteractiveRequest.driverMemory.foreach(builder.driverMemory)
     createInteractiveRequest.executorCores.foreach(builder.executorCores)
+    createInteractiveRequest.executorMemory.foreach(builder.executorMemory)
+    createInteractiveRequest.files.map(RelativePath).foreach(builder.file)
+    createInteractiveRequest.jars.map(RelativePath).foreach(builder.jar)
+    createInteractiveRequest.proxyUser.foreach(builder.proxyUser)
+    createInteractiveRequest.pyFiles.map(RelativePath).foreach(builder.pyFile)
 
     builder.redirectOutput(Redirect.PIPE)
     builder.redirectErrorStream(redirect = true)