Browse Source

[livy] Create a SparkProcessBuilderFactory

Erick Tryzelaar 10 years ago
parent
commit
539d0ad
17 changed files with 228 additions and 213 deletions
  1. 15 0
      apps/spark/java/livy-core/src/main/scala/com/cloudera/hue/livy/spark/SparkProcessBuilderFactory.scala
  2. 14 2
      apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/Main.scala
  3. 2 1
      apps/spark/java/livy-server/src/test/scala/com/cloudera/hue/livy/server/batch/BatchServletSpec.scala
  4. 9 3
      apps/spark/java/livy-server/src/test/scala/com/cloudera/hue/livy/server/interactive/InteractiveSessionServletSpec.scala
  5. 36 3
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionFactory.scala
  6. 3 32
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionProcess.scala
  7. 6 5
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionProcessFactory.scala
  8. 3 35
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionYarn.scala
  9. 10 5
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionYarnFactory.scala
  10. 61 3
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionFactory.scala
  11. 9 51
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionProcess.scala
  12. 30 8
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionProcessFactory.scala
  13. 9 52
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionYarn.scala
  14. 7 6
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionYarnFactory.scala
  15. 3 5
      apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveWebSession.scala
  16. 5 1
      apps/spark/java/livy-spark/src/test/scala/com/cloudera/hue/livy/spark/batch/BatchProcessSpec.scala
  17. 6 1
      apps/spark/java/livy-spark/src/test/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionProcessSpec.scala

+ 15 - 0
apps/spark/java/livy-core/src/main/scala/com/cloudera/hue/livy/spark/SparkProcessBuilderFactory.scala

@@ -0,0 +1,15 @@
+package com.cloudera.hue.livy.spark
+
+import com.cloudera.hue.livy.LivyConf
+
+object SparkProcessBuilderFactory {
+  def apply(livyConf: LivyConf): SparkProcessBuilderFactory = {
+    new SparkProcessBuilderFactory(livyConf)
+  }
+}
+
+class SparkProcessBuilderFactory(val livyConf: LivyConf) {
+  def builder() = {
+    SparkProcessBuilder(livyConf)
+  }
+}

+ 14 - 2
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/Main.scala

@@ -26,8 +26,10 @@ import com.cloudera.hue.livy.server.batch.BatchSessionServlet
 import com.cloudera.hue.livy.server.interactive.InteractiveSessionServlet
 import com.cloudera.hue.livy.sessions.batch.BatchSession
 import com.cloudera.hue.livy.sessions.interactive.InteractiveSession
+import com.cloudera.hue.livy.spark.SparkProcessBuilderFactory
 import com.cloudera.hue.livy.spark.batch.{BatchSessionProcessFactory, BatchSessionYarnFactory}
 import com.cloudera.hue.livy.spark.interactive.{InteractiveSessionYarnFactory, InteractiveSessionProcessFactory}
+import com.cloudera.hue.livy.yarn.Client
 import org.scalatra._
 import org.scalatra.metrics.MetricsBootstrap
 import org.scalatra.metrics.MetricsSupportExtensions._
@@ -162,11 +164,21 @@ class ScalatraBootstrap
 
       info(f"Using $sessionFactoryKind sessions")
 
+      val processFactory = new SparkProcessBuilderFactory(livyConf)
+
       val (sessionFactory, batchFactory) = sessionFactoryKind match {
         case LivyConf.Process() =>
-          (new InteractiveSessionProcessFactory(livyConf), new BatchSessionProcessFactory(livyConf))
+          val interactiveFactory = new InteractiveSessionProcessFactory(processFactory)
+          val batchFactory = new BatchSessionProcessFactory(processFactory)
+
+          (interactiveFactory, batchFactory)
+
         case LivyConf.Yarn() =>
-          (new InteractiveSessionYarnFactory(livyConf), new BatchSessionYarnFactory(livyConf))
+          val client = new Client(livyConf)
+          val interactiveFactory = new InteractiveSessionYarnFactory(client, processFactory)
+          val batchFactory = new BatchSessionYarnFactory(client, processFactory)
+
+          (interactiveFactory, batchFactory)
       }
 
       sessionManager = new SessionManager(livyConf, sessionFactory)

+ 2 - 1
apps/spark/java/livy-server/src/test/scala/com/cloudera/hue/livy/server/batch/BatchServletSpec.scala

@@ -24,6 +24,7 @@ import java.util.concurrent.TimeUnit
 
 import com.cloudera.hue.livy.server.SessionManager
 import com.cloudera.hue.livy.sessions.SessionState
+import com.cloudera.hue.livy.spark.SparkProcessBuilderFactory
 import com.cloudera.hue.livy.spark.batch.{BatchSessionProcessFactory, CreateBatchRequest}
 import com.cloudera.hue.livy.{LivyConf, Utils}
 import org.json4s.JsonAST.{JArray, JInt, JObject, JString}
@@ -55,7 +56,7 @@ class BatchServletSpec extends ScalatraSuite with FunSpecLike with BeforeAndAfte
   }
 
   val livyConf = new LivyConf()
-  val batchFactory = new BatchSessionProcessFactory(livyConf)
+  val batchFactory = new BatchSessionProcessFactory(SparkProcessBuilderFactory(livyConf))
   val batchManager = new SessionManager(livyConf, batchFactory)
   val servlet = new BatchSessionServlet(batchManager)
 

+ 9 - 3
apps/spark/java/livy-server/src/test/scala/com/cloudera/hue/livy/server/interactive/InteractiveSessionServletSpec.scala

@@ -26,6 +26,7 @@ import com.cloudera.hue.livy.msgs.ExecuteRequest
 import com.cloudera.hue.livy.server.SessionManager
 import com.cloudera.hue.livy.sessions._
 import com.cloudera.hue.livy.sessions.interactive.{InteractiveSession, Statement}
+import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilderFactory}
 import com.cloudera.hue.livy.spark.interactive.{CreateInteractiveRequest, InteractiveSessionFactory}
 import org.json4s.JsonAST.{JArray, JInt, JObject, JString}
 import org.json4s.jackson.JsonMethods._
@@ -77,14 +78,19 @@ class InteractiveSessionServletSpec extends ScalatraSuite with FunSpecLike {
     override def interrupt(): Future[Unit] = ???
   }
 
-  class MockInteractiveSessionFactory() extends InteractiveSessionFactory {
-    override def create(id: Int, createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
+  class MockInteractiveSessionFactory(processFactory: SparkProcessBuilderFactory)
+    extends InteractiveSessionFactory(processFactory) {
+
+    protected override def create(id: Int,
+                                  process: SparkProcess,
+                                  request: CreateInteractiveRequest): InteractiveSession = {
       new MockInteractiveSession(id)
     }
   }
 
   val livyConf = new LivyConf()
-  val sessionManager = new SessionManager(livyConf, new MockInteractiveSessionFactory())
+  val processFactory = new SparkProcessBuilderFactory(livyConf)
+  val sessionManager = new SessionManager(livyConf, new MockInteractiveSessionFactory(processFactory))
   val servlet = new InteractiveSessionServlet(sessionManager)
 
   addServlet(servlet, "/*")

+ 36 - 3
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionFactory.scala

@@ -18,13 +18,46 @@
 
 package com.cloudera.hue.livy.spark.batch
 
+import java.lang.ProcessBuilder.Redirect
+
 import com.cloudera.hue.livy.sessions.SessionFactory
 import com.cloudera.hue.livy.sessions.batch.BatchSession
+import com.cloudera.hue.livy.spark.SparkProcessBuilder.RelativePath
+import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilder, SparkProcessBuilderFactory}
 import org.json4s.JValue
 
-abstract class BatchSessionFactory extends SessionFactory[BatchSession] {
+abstract class BatchSessionFactory(factory: SparkProcessBuilderFactory) extends SessionFactory[BatchSession] {
   override def create(id: Int, createRequest: JValue) =
     create(id, createRequest.extract[CreateBatchRequest])
 
-  def create(id: Int, createRequest: CreateBatchRequest): BatchSession
-}
+  def create(id: Int, request: CreateBatchRequest): BatchSession = {
+    val builder = sparkBuilder(request)
+    val process = builder.start(RelativePath(request.file), request.args)
+    create(id, process)
+  }
+
+  protected def create(id: Int, process: SparkProcess): BatchSession
+
+  protected def sparkBuilder(request: CreateBatchRequest): SparkProcessBuilder = {
+    val builder = factory.builder()
+
+    request.proxyUser.foreach(builder.proxyUser)
+    request.className.foreach(builder.className)
+    request.jars.map(RelativePath).foreach(builder.jar)
+    request.pyFiles.map(RelativePath).foreach(builder.pyFile)
+    request.files.map(RelativePath).foreach(builder.file)
+    request.driverMemory.foreach(builder.driverMemory)
+    request.driverCores.foreach(builder.driverCores)
+    request.executorMemory.foreach(builder.executorMemory)
+    request.executorCores.foreach(builder.executorCores)
+    request.numExecutors.foreach(builder.numExecutors)
+    request.archives.map(RelativePath).foreach(builder.archive)
+    request.queue.foreach(builder.queue)
+    request.name.foreach(builder.name)
+
+    builder.redirectOutput(Redirect.PIPE)
+    builder.redirectErrorStream(true)
+
+    builder
+  }
+}

+ 3 - 32
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionProcess.scala

@@ -18,46 +18,17 @@
 
 package com.cloudera.hue.livy.spark.batch
 
-import java.lang.ProcessBuilder.Redirect
-
+import com.cloudera.hue.livy.LineBufferedProcess
 import com.cloudera.hue.livy.sessions.SessionState
 import com.cloudera.hue.livy.sessions.batch.BatchSession
-import com.cloudera.hue.livy.spark.SparkProcessBuilder
-import com.cloudera.hue.livy.spark.SparkProcessBuilder.RelativePath
-import com.cloudera.hue.livy.{LineBufferedProcess, LivyConf}
+import com.cloudera.hue.livy.spark.SparkProcess
 
 import scala.concurrent.{ExecutionContext, ExecutionContextExecutor, Future}
 
 object BatchSessionProcess {
-  def apply(livyConf: LivyConf, id: Int, createBatchRequest: CreateBatchRequest): BatchSession = {
-    val builder = sparkBuilder(livyConf, createBatchRequest)
-
-    val process = builder.start(RelativePath(createBatchRequest.file), createBatchRequest.args)
+  def apply(id: Int, process: SparkProcess): BatchSession = {
     new BatchSessionProcess(id, process)
   }
-
-  private def sparkBuilder(livyConf: LivyConf, createBatchRequest: CreateBatchRequest): SparkProcessBuilder = {
-    val builder = SparkProcessBuilder(livyConf)
-
-    createBatchRequest.className.foreach(builder.className)
-    createBatchRequest.jars.map(RelativePath).foreach(builder.jar)
-    createBatchRequest.pyFiles.map(RelativePath).foreach(builder.pyFile)
-    createBatchRequest.files.map(RelativePath).foreach(builder.file)
-    createBatchRequest.driverMemory.foreach(builder.driverMemory)
-    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)
-    createBatchRequest.proxyUser.foreach(builder.proxyUser)
-    createBatchRequest.queue.foreach(builder.queue)
-    createBatchRequest.name.foreach(builder.name)
-
-    builder.redirectOutput(Redirect.PIPE)
-    builder.redirectErrorStream(true)
-
-    builder
-  }
 }
 
 private class BatchSessionProcess(val id: Int,

+ 6 - 5
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionProcessFactory.scala

@@ -18,12 +18,13 @@
 
 package com.cloudera.hue.livy.spark.batch
 
-import com.cloudera.hue.livy.LivyConf
 import com.cloudera.hue.livy.sessions.batch.BatchSession
+import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilderFactory}
 
-class BatchSessionProcessFactory(livyConf: LivyConf)
-  extends BatchSessionFactory
+class BatchSessionProcessFactory(processFactory: SparkProcessBuilderFactory)
+  extends BatchSessionFactory(processFactory)
 {
-  override def create(id: Int, createBatchRequest: CreateBatchRequest): BatchSession =
-    BatchSessionProcess(livyConf, id, createBatchRequest)
+  protected override def create(id: Int, process: SparkProcess): BatchSession = {
+    BatchSessionProcess(id, process)
+  }
 }

+ 3 - 35
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionYarn.scala

@@ -18,56 +18,24 @@
 
 package com.cloudera.hue.livy.spark.batch
 
-import java.lang.ProcessBuilder.Redirect
-
+import com.cloudera.hue.livy.LineBufferedProcess
 import com.cloudera.hue.livy.sessions._
 import com.cloudera.hue.livy.sessions.batch.BatchSession
-import com.cloudera.hue.livy.spark.SparkProcessBuilder
-import com.cloudera.hue.livy.spark.SparkProcessBuilder.RelativePath
+import com.cloudera.hue.livy.spark.SparkProcess
 import com.cloudera.hue.livy.yarn._
-import com.cloudera.hue.livy.{LineBufferedProcess, LivyConf}
 
 import scala.annotation.tailrec
 import scala.concurrent.{ExecutionContext, ExecutionContextExecutor, Future}
 
 object BatchSessionYarn {
-
   implicit def executor: ExecutionContextExecutor = ExecutionContext.global
 
-  def apply(livyConf: LivyConf, client: Client, id: Int, createBatchRequest: CreateBatchRequest): BatchSession = {
-    val builder = sparkBuilder(livyConf, createBatchRequest)
-
-    val process = builder.start(RelativePath(createBatchRequest.file), createBatchRequest.args)
+  def apply(client: Client, id: Int, process: SparkProcess): BatchSession = {
     val job = Future {
       client.getJobFromProcess(process)
     }
     new BatchSessionYarn(id, process, job)
   }
-
-  private def sparkBuilder(livyConf: LivyConf, createBatchRequest: CreateBatchRequest): SparkProcessBuilder = {
-    val builder = SparkProcessBuilder(livyConf)
-
-    builder.master("yarn-cluster")
-
-    createBatchRequest.proxyUser.foreach(builder.proxyUser)
-    createBatchRequest.className.foreach(builder.className)
-    createBatchRequest.jars.map(RelativePath).foreach(builder.jar)
-    createBatchRequest.pyFiles.map(RelativePath).foreach(builder.pyFile)
-    createBatchRequest.files.map(RelativePath).foreach(builder.file)
-    createBatchRequest.driverMemory.foreach(builder.driverMemory)
-    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)
-    createBatchRequest.queue.foreach(builder.queue)
-    createBatchRequest.name.foreach(builder.name)
-
-    builder.redirectOutput(Redirect.PIPE)
-    builder.redirectErrorStream(true)
-
-    builder
-  }
 }
 
 private class BatchSessionYarn(val id: Int, process: LineBufferedProcess, jobFuture: Future[Job]) extends BatchSession {

+ 10 - 5
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/batch/BatchSessionYarnFactory.scala

@@ -19,13 +19,18 @@
 package com.cloudera.hue.livy.spark.batch
 
 import com.cloudera.hue.livy.LivyConf
-import com.cloudera.hue.livy.sessions.batch.BatchSession
+import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilderFactory}
 import com.cloudera.hue.livy.yarn.Client
 
-class BatchSessionYarnFactory(livyConf: LivyConf) extends BatchSessionFactory {
+class BatchSessionYarnFactory(client: Client, factory: SparkProcessBuilderFactory)
+  extends BatchSessionFactory(factory) {
 
-  val client = new Client(livyConf)
+  protected override def create(id: Int, process: SparkProcess) =
+    BatchSessionYarn(client, id, process)
 
-  def create(id: Int, createBatchRequest: CreateBatchRequest): BatchSession =
-    BatchSessionYarn(livyConf, client, id, createBatchRequest)
+  override def sparkBuilder(request: CreateBatchRequest) = {
+    val builder = super.sparkBuilder(request)
+    builder.master("yarn-cluster")
+    builder
+  }
 }

+ 61 - 3
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionFactory.scala

@@ -18,16 +18,74 @@
 
 package com.cloudera.hue.livy.spark.interactive
 
+import java.lang.ProcessBuilder.Redirect
+
 import com.cloudera.hue.livy.sessions.interactive.InteractiveSession
-import com.cloudera.hue.livy.sessions.{SessionFactory, SessionKindSerializer}
+import com.cloudera.hue.livy.sessions.{PySpark, SessionFactory, SessionKindSerializer}
+import com.cloudera.hue.livy.spark.SparkProcessBuilder.{AbsolutePath, RelativePath}
+import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilder, SparkProcessBuilderFactory}
+import com.cloudera.hue.livy.{LivyConf, Utils}
 import org.json4s.{DefaultFormats, Formats, JValue}
 
-trait InteractiveSessionFactory extends SessionFactory[InteractiveSession] {
+object InteractiveSessionFactory {
+  private val CONF_LIVY_JAR = "livy.repl.jar"
+}
+
+abstract class InteractiveSessionFactory(processFactory: SparkProcessBuilderFactory)
+  extends SessionFactory[InteractiveSession] {
 
   override protected implicit def jsonFormats: Formats = DefaultFormats ++ List(SessionKindSerializer)
 
   override def create(id: Int, createRequest: JValue) =
     create(id, createRequest.extract[CreateInteractiveRequest])
 
-  def create(id: Int, createRequest: CreateInteractiveRequest): InteractiveSession
+  def create(id: Int, request: CreateInteractiveRequest): InteractiveSession = {
+    val builder = sparkBuilder(id, request)
+    val kind = request.kind.toString
+    val process = builder.start(AbsolutePath(livyJar(processFactory.livyConf)), List(kind))
+
+    create(id, process, request)
+  }
+
+  protected def create(id: Int, process: SparkProcess, request: CreateInteractiveRequest): InteractiveSession
+
+  protected def sparkBuilder(id: Int, request: CreateInteractiveRequest): SparkProcessBuilder = {
+    val builder = processFactory.builder()
+
+    builder.className("com.cloudera.hue.livy.repl.Main")
+    request.archives.map(RelativePath).foreach(builder.archive)
+    request.driverCores.foreach(builder.driverCores)
+    request.driverMemory.foreach(builder.driverMemory)
+    request.executorCores.foreach(builder.executorCores)
+    request.executorMemory.foreach(builder.executorMemory)
+    request.numExecutors.foreach(builder.numExecutors)
+    request.files.map(RelativePath).foreach(builder.file)
+    request.jars.map(RelativePath).foreach(builder.jar)
+    request.proxyUser.foreach(builder.proxyUser)
+    request.pyFiles.map(RelativePath).foreach(builder.pyFile)
+    request.queue.foreach(builder.queue)
+    request.name.foreach(builder.name)
+
+    val callbackUrl = System.getProperty("livy.server.callback-url")
+    val url = f"$callbackUrl/sessions/$id/callback"
+
+    builder.driverJavaOptions(f"-Dlivy.repl.callback-url=$url -Dlivy.repl.port=0")
+
+    request.kind match {
+      case PySpark() => builder.conf("spark.yarn.isPython", "true")
+      case _ =>
+    }
+
+    builder.env("LIVY_PORT", "0")
+
+    builder.redirectOutput(Redirect.PIPE)
+    builder.redirectErrorStream(true)
+
+    builder
+  }
+
+  private def livyJar(livyConf: LivyConf) = {
+    livyConf.getOption(InteractiveSessionFactory.CONF_LIVY_JAR)
+      .getOrElse(Utils.jarOfClass(getClass).head)
+  }
 }

+ 9 - 51
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionProcess.scala

@@ -18,14 +18,12 @@
 
 package com.cloudera.hue.livy.spark.interactive
 
-import java.lang.ProcessBuilder.Redirect
 import java.net.URL
 
-import com.cloudera.hue.livy.sessions._
+import com.cloudera.hue.livy.Logging
+import com.cloudera.hue.livy.sessions.SessionState
 import com.cloudera.hue.livy.sessions.interactive.InteractiveSession
-import com.cloudera.hue.livy.spark.SparkProcessBuilder.{AbsolutePath, RelativePath}
-import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilder}
-import com.cloudera.hue.livy.{LivyConf, Logging, Utils}
+import com.cloudera.hue.livy.spark.SparkProcess
 
 import scala.annotation.tailrec
 import scala.concurrent.Future
@@ -36,57 +34,17 @@ object InteractiveSessionProcess extends Logging {
   val CONF_LIVY_REPL_CALLBACK_URL = "livy.repl.callback-url"
   val CONF_LIVY_REPL_DRIVER_CLASS_PATH = "livy.repl.driverClassPath"
 
-  def apply(livyConf: LivyConf,
-            id: Int,
+  def apply(id: Int,
+            process: SparkProcess,
             createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
-    val process = startProcess(livyConf, id, createInteractiveRequest)
-    new InteractiveSessionProcess(id, createInteractiveRequest, process)
-  }
-
-  // Loop until we've started a process with a valid port.
-  private def startProcess(livyConf: LivyConf, id: Int, createInteractiveRequest: CreateInteractiveRequest): SparkProcess = {
-
-    val builder = new SparkProcessBuilder(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)
-    createInteractiveRequest.numExecutors.foreach(builder.numExecutors)
-    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)
-    createInteractiveRequest.queue.foreach(builder.queue)
-    createInteractiveRequest.name.foreach(builder.name)
-
-    sys.env.get("LIVY_REPL_JAVA_OPTS").foreach(builder.driverJavaOptions)
-    livyConf.getOption(CONF_LIVY_REPL_DRIVER_CLASS_PATH).foreach(builder.driverClassPath)
-
-    livyConf.getOption(CONF_LIVY_REPL_CALLBACK_URL).foreach { case callbackUrl =>
-      builder.env("LIVY_CALLBACK_URL", f"$callbackUrl/sessions/$id/callback")
-    }
-
-    builder.env("LIVY_PORT", "0")
-
-    builder.redirectOutput(Redirect.PIPE)
-    builder.redirectErrorStream(true)
-
-    builder.start(AbsolutePath(livyJar(livyConf)), List(createInteractiveRequest.kind.toString))
-  }
-
-  private def livyJar(conf: LivyConf): String = {
-    conf.getOption(CONF_LIVY_REPL_JAR).getOrElse {
-      Utils.jarOfClass(getClass).head
-    }
+    new InteractiveSessionProcess(id, process, createInteractiveRequest)
   }
 }
 
 private class InteractiveSessionProcess(id: Int,
-                                        createInteractiveRequest: CreateInteractiveRequest,
-                                        process: SparkProcess) extends InteractiveWebSession(id, createInteractiveRequest) {
+                                        process: SparkProcess,
+                                        createInteractiveRequest: CreateInteractiveRequest)
+  extends InteractiveWebSession(id, createInteractiveRequest) {
 
   val stdoutThread = new Thread {
     override def run() = {

+ 30 - 8
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionProcessFactory.scala

@@ -18,17 +18,39 @@
 
 package com.cloudera.hue.livy.spark.interactive
 
-import com.cloudera.hue.livy.LivyConf
 import com.cloudera.hue.livy.sessions.interactive.InteractiveSession
+import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilderFactory}
 
 import scala.concurrent.ExecutionContext
 
-class InteractiveSessionProcessFactory(livyConf: LivyConf) extends InteractiveSessionFactory {
+object InteractiveSessionProcessFactory {
+  val CONF_LIVY_REPL_JAR = "livy.repl.jar"
+  val CONF_LIVY_REPL_CALLBACK_URL = "livy.repl.callback-url"
+  val CONF_LIVY_REPL_DRIVER_CLASS_PATH = "livy.repl.driverClassPath"
+}
 
-   implicit def executor: ExecutionContext = ExecutionContext.global
+class InteractiveSessionProcessFactory(processFactory: SparkProcessBuilderFactory)
+  extends InteractiveSessionFactory(processFactory) {
 
-   override def create(id: Int,
-                       createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
-     InteractiveSessionProcess(livyConf, id, createInteractiveRequest)
-   }
- }
+  implicit def executor: ExecutionContext = ExecutionContext.global
+
+  protected override def create(id: Int, process: SparkProcess, createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
+    InteractiveSessionProcess(id, process, createInteractiveRequest)
+  }
+
+  override def sparkBuilder(id: Int, request: CreateInteractiveRequest) = {
+    val builder = super.sparkBuilder(id, request)
+
+    sys.env.get("LIVY_REPL_JAVA_OPTS").foreach(builder.driverJavaOptions)
+    processFactory.livyConf.getOption(InteractiveSessionProcessFactory.CONF_LIVY_REPL_DRIVER_CLASS_PATH)
+      .foreach(builder.driverClassPath)
+
+    processFactory.livyConf.getOption(InteractiveSessionProcessFactory.CONF_LIVY_REPL_CALLBACK_URL)
+      .foreach { case callbackUrl =>
+        builder.env("LIVY_CALLBACK_URL", f"$callbackUrl/sessions/$id/callback")
+      }
+
+    builder.env("LIVY_PORT", "0")
+    builder
+  }
+}

+ 9 - 52
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionYarn.scala

@@ -18,15 +18,12 @@
 
 package com.cloudera.hue.livy.spark.interactive
 
-import java.lang.ProcessBuilder.Redirect
 import java.util.concurrent.TimeUnit
 
+import com.cloudera.hue.livy.sessions.SessionState
 import com.cloudera.hue.livy.sessions.interactive.InteractiveSession
-import com.cloudera.hue.livy.sessions.{PySpark, SessionState}
-import com.cloudera.hue.livy.spark.SparkProcessBuilder.{AbsolutePath, RelativePath}
-import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilder}
+import com.cloudera.hue.livy.spark.SparkProcess
 import com.cloudera.hue.livy.yarn.Client
-import com.cloudera.hue.livy.{LivyConf, Utils}
 
 import scala.concurrent.duration._
 import scala.concurrent.{Await, ExecutionContext, ExecutionContextExecutor, Future}
@@ -34,61 +31,21 @@ import scala.concurrent.{Await, ExecutionContext, ExecutionContextExecutor, Futu
 object InteractiveSessionYarn {
   protected implicit def executor: ExecutionContextExecutor = ExecutionContext.global
 
-  private val CONF_LIVY_JAR = "livy.yarn.jar"
+  private lazy val regex = """Application report for (\w+)""".r.unanchored
 
-  def apply(livyConf: LivyConf,
-            client: Client,
+  def apply(client: Client,
             id: Int,
-            createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
-    val callbackUrl = System.getProperty("livy.server.callback-url")
-    val url = f"$callbackUrl/sessions/$id/callback"
-
-    val builder = SparkProcessBuilder(livyConf)
-
-    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.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)
-    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)
-    createInteractiveRequest.queue.foreach(builder.queue)
-    createInteractiveRequest.name.foreach(builder.name)
-
-    val kind = createInteractiveRequest.kind.toString
-
-    createInteractiveRequest.kind match {
-      case PySpark() => builder.conf("spark.yarn.isPython", "true")
-      case _ =>
-    }
-
-    builder.redirectOutput(Redirect.PIPE)
-    builder.redirectErrorStream(true)
-
-    val process = builder.start(AbsolutePath(livyJar(livyConf)), List(kind))
-
-    new InteractiveSessionYarn(id, client, process, createInteractiveRequest)
-  }
-
-  private def livyJar(livyConf: LivyConf) = {
-    if (livyConf.contains(CONF_LIVY_JAR)) {
-      livyConf.get(CONF_LIVY_JAR)
-    } else {
-      Utils.jarOfClass(classOf[Client]).head
-    }
+            process: SparkProcess,
+            request: CreateInteractiveRequest): InteractiveSession = {
+    new InteractiveSessionYarn(id, client, process, request)
   }
 }
 
 private class InteractiveSessionYarn(id: Int,
                                      client: Client,
                                      process: SparkProcess,
-                                     createInteractiveRequest: CreateInteractiveRequest)
-  extends InteractiveWebSession(id, createInteractiveRequest) {
+                                     request: CreateInteractiveRequest)
+  extends InteractiveWebSession(id, request) {
 
   // Error out the job if the process errors out.
   Future {

+ 7 - 6
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionYarnFactory.scala

@@ -18,20 +18,21 @@
 
 package com.cloudera.hue.livy.spark.interactive
 
-import com.cloudera.hue.livy.LivyConf
 import com.cloudera.hue.livy.sessions.interactive.InteractiveSession
+import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilderFactory}
 import com.cloudera.hue.livy.yarn.Client
 
 import scala.concurrent.ExecutionContext
 
-class InteractiveSessionYarnFactory(livyConf: LivyConf) extends InteractiveSessionFactory {
+class InteractiveSessionYarnFactory(client: Client, processFactory: SparkProcessBuilderFactory)
+  extends InteractiveSessionFactory(processFactory) {
 
    implicit def executor: ExecutionContext = ExecutionContext.global
 
-   val client = new Client(livyConf)
-
-   override def create(id: Int, createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
-     InteractiveSessionYarn(livyConf, client, id, createInteractiveRequest)
+   protected  override def create(id: Int,
+                                  process: SparkProcess,
+                                  request: CreateInteractiveRequest): InteractiveSession = {
+     InteractiveSessionYarn(client, id, process, request)
    }
 
    override def close(): Unit = {

+ 3 - 5
apps/spark/java/livy-spark/src/main/scala/com/cloudera/hue/livy/spark/interactive/InteractiveWebSession.scala

@@ -34,9 +34,7 @@ import scala.annotation.tailrec
 import scala.concurrent.duration.Duration
 import scala.concurrent.{Future, _}
 
-abstract class InteractiveWebSession(val id: Int,
-                                     createInteractiveRequest: CreateInteractiveRequest)
-  extends InteractiveSession with Logging {
+abstract class InteractiveWebSession(val id: Int, request: CreateInteractiveRequest) extends InteractiveSession with Logging {
 
   protected implicit def executor: ExecutionContextExecutor = ExecutionContext.global
   protected implicit def jsonFormats: Formats = DefaultFormats
@@ -49,9 +47,9 @@ abstract class InteractiveWebSession(val id: Int,
   private[this] var _executedStatements = 0
   private[this] var _statements = IndexedSeq[Statement]()
 
-  override def kind = createInteractiveRequest.kind
+  override def kind = request.kind
 
-  override def proxyUser = createInteractiveRequest.proxyUser
+  override def proxyUser = request.proxyUser
 
   override def url: Option[URL] = _url
 

+ 5 - 1
apps/spark/java/livy-spark/src/test/scala/com/cloudera/hue/livy/spark/batch/BatchProcessSpec.scala

@@ -23,6 +23,7 @@ import java.nio.file.{Files, Path}
 import java.util.concurrent.TimeUnit
 
 import com.cloudera.hue.livy.sessions.SessionState
+import com.cloudera.hue.livy.spark.SparkProcessBuilderFactory
 import com.cloudera.hue.livy.{LivyConf, Utils}
 import org.scalatest.{BeforeAndAfterAll, FunSpec, ShouldMatchers}
 
@@ -53,7 +54,10 @@ class BatchProcessSpec
       val req = CreateBatchRequest(
         file = script.toString
       )
-      val batch = BatchSessionProcess(new LivyConf(), 0, req)
+
+      val livyConf = new LivyConf()
+      val builder = new BatchSessionProcessFactory(SparkProcessBuilderFactory(livyConf))
+      val batch = builder.create(0, req)
 
       Utils.waitUntil({ () => !batch.state.isActive }, Duration(10, TimeUnit.SECONDS))
       (batch.state match {

+ 6 - 1
apps/spark/java/livy-spark/src/test/scala/com/cloudera/hue/livy/spark/interactive/InteractiveSessionProcessSpec.scala

@@ -20,6 +20,7 @@ package com.cloudera.hue.livy.spark.interactive
 
 import com.cloudera.hue.livy.LivyConf
 import com.cloudera.hue.livy.sessions.{BaseInteractiveSessionSpec, PySpark}
+import com.cloudera.hue.livy.spark.SparkProcessBuilderFactory
 import org.scalatest.{BeforeAndAfter, FunSpecLike, Matchers}
 
 class InteractiveSessionProcessSpec
@@ -31,5 +32,9 @@ class InteractiveSessionProcessSpec
   val livyConf = new LivyConf()
   livyConf.set("livy.repl.driverClassPath", sys.props("java.class.path"))
 
-  def createSession() = InteractiveSessionProcess(livyConf, 0, CreateInteractiveRequest(kind = PySpark()))
+  def createSession() = {
+    val processFactory = new SparkProcessBuilderFactory(livyConf)
+    val interactiveFactory = new InteractiveSessionProcessFactory(processFactory)
+    interactiveFactory.create(0, CreateInteractiveRequest(kind = PySpark()))
+  }
 }