Explorar o código

HUE-3046 [livy] Pass the callback url and port through a config variable

This allows ./bin/livy-repl to keep working without needing
a dummy callback url to be defined.
Erick Tryzelaar %!s(int64=10) %!d(string=hai) anos
pai
achega
af01a43

+ 6 - 13
apps/spark/java/livy-repl/src/main/scala/com/cloudera/hue/livy/repl/Main.scala

@@ -37,22 +37,17 @@ import _root_.scala.concurrent.duration._
 import _root_.scala.concurrent.{Await, ExecutionContext}
 
 object Main extends Logging {
-
   val SESSION_KIND = "livy.repl.session.kind"
+  val CALLBACK_URL = "livy.repl.callbackUrl"
   val PYSPARK_SESSION = "pyspark"
   val SPARK_SESSION = "spark"
   val SPARKR_SESSION = "sparkr"
 
   def main(args: Array[String]): Unit = {
 
-    val host = Option(System.getProperty("livy.repl.host"))
-      .orElse(sys.env.get("LIVY_HOST"))
-      .getOrElse("0.0.0.0")
-
-    val port = Option(System.getProperty("livy.repl.port"))
-      .orElse(sys.env.get("LIVY_PORT"))
-      .getOrElse("8999").toInt
-
+    val host = sys.props.getOrElse("spark.livy.host", "0.0.0.0")
+    val port = sys.props.getOrElse("spark.livy.port", "8999").toInt
+    val callbackUrl = sys.props.get("spark.livy.callbackUrl")
 
     if (args.length != 1) {
       println("Must specify either `pyspark`/`spark`/`sparkr` for the session kind")
@@ -74,6 +69,7 @@ object Main extends Logging {
     server.context.addEventListener(new ScalatraListener)
     server.context.setInitParameter(ScalatraListener.LifeCycleKey, classOf[ScalatraBootstrap].getCanonicalName)
     server.context.setInitParameter(SESSION_KIND, session_kind)
+    callbackUrl.foreach(server.context.setInitParameter(CALLBACK_URL, _))
 
     server.start()
 
@@ -112,11 +108,8 @@ class ScalatraBootstrap extends LifeCycle with Logging {
 
       context.mount(new WebApp(session), "/*")
 
-      val callbackUrl = Option(System.getProperty("livy.repl.callback-url"))
-        .orElse(sys.env.get("LIVY_CALLBACK_URL"))
-
       // See if we want to notify someone that we've started on a url
-      callbackUrl.foreach(notifyCallback)
+      Option(context.getInitParameter(Main.CALLBACK_URL)).foreach(notifyCallback)
     } catch {
       case e: Throwable =>
         println(f"Exception thrown when initializing server: $e")

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

@@ -58,7 +58,7 @@ object Main {
     server.start()
 
     try {
-      System.setProperty("livy.server.callback-url", f"http://${server.host}:${server.port}")
+      System.setProperty("livy.server.serverUrl", f"http://${server.host}:${server.port}")
     } finally {
       server.join()
       server.stop()

+ 2 - 0
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/interactive/InteractiveSessionServlet.scala

@@ -51,6 +51,8 @@ class InteractiveSessionServlet(sessionManager: SessionManager[InteractiveSessio
         if (session.state == SessionState.Starting()) {
           session.url = new URL(callback.url)
           Accepted()
+        } else if (session.state.isActive) {
+          Ok()
         } else {
           BadRequest("Session is in wrong state")
         }

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

@@ -28,12 +28,20 @@ import com.cloudera.hue.livy.{LivyConf, Utils}
 import org.json4s.{DefaultFormats, Formats, JValue}
 
 object InteractiveSessionFactory {
-  private val CONF_LIVY_JAR = "livy.repl.jar"
+  private val LivyReplDriverClassPath = "livy.repl.driverClassPath"
+  private val LivyReplJar = "livy.repl.jar"
+  private val LivyServerUrl = "livy.server.serverUrl"
+  private val SparkDriverExtraJavaOptions = "spark.driver.extraDriverOptions"
+  private val SparkLivyCallbackUrl = "spark.livy.callbackUrl"
+  private val SparkLivyPort = "spark.livy.port"
+  private val SparkYarnIsPython = "spark.yarn.isPython"
 }
 
 abstract class InteractiveSessionFactory(processFactory: SparkProcessBuilderFactory)
   extends SessionFactory[InteractiveSession] {
 
+  import InteractiveSessionFactory._
+
   override protected implicit def jsonFormats: Formats = DefaultFormats ++ List(SessionKindSerializer)
 
   override def create(id: Int, createRequest: JValue) =
@@ -67,23 +75,20 @@ abstract class InteractiveSessionFactory(processFactory: SparkProcessBuilderFact
     request.queue.foreach(builder.queue)
     request.name.foreach(builder.name)
 
-    val callbackUrl = System.getProperty("livy.server.callback-url")
-    val driverOptions =
-      f"-Dlivy.repl.callback-url=$callbackUrl/sessions/$id/callback " +
-      "-Dlivy.repl.port=0"
-
-    builder.conf(
-      "spark.driver.extraJavaOptions",
-      builder.conf("spark.driver.extraJavaOptions")
-        .getOrElse("") + driverOptions,
-        admin = true)
-
     request.kind match {
-      case PySpark() => builder.conf("spark.yarn.isPython", "true", admin = true)
+      case PySpark() => builder.conf(SparkYarnIsPython, "true", admin = true)
       case _ =>
     }
 
-    builder.env("LIVY_PORT", "0")
+    processFactory.livyConf.getOption(LivyReplDriverClassPath)
+      .foreach(builder.driverClassPath)
+
+    sys.props.get(LivyServerUrl).foreach { serverUrl =>
+      val callbackUrl = f"$serverUrl/sessions/$id/callback"
+      builder.conf(SparkLivyCallbackUrl, callbackUrl, admin = true)
+    }
+
+    builder.conf(SparkLivyPort, "0", admin = true)
 
     builder.redirectOutput(Redirect.PIPE)
     builder.redirectErrorStream(true)
@@ -92,7 +97,7 @@ abstract class InteractiveSessionFactory(processFactory: SparkProcessBuilderFact
   }
 
   private def livyJar(livyConf: LivyConf) = {
-    livyConf.getOption(InteractiveSessionFactory.CONF_LIVY_JAR)
+    livyConf.getOption(LivyReplJar)
       .getOrElse(Utils.jarOfClass(getClass).head)
   }
 }

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

@@ -21,19 +21,13 @@ package com.cloudera.hue.livy.spark.interactive
 import java.net.URL
 
 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.SparkProcess
 
 import scala.annotation.tailrec
-import scala.concurrent.Future
 
 object InteractiveSessionProcess extends Logging {
 
-  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"
-
   def apply(id: Int,
             process: SparkProcess,
             createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
@@ -43,8 +37,8 @@ object InteractiveSessionProcess extends Logging {
 
 private class InteractiveSessionProcess(id: Int,
                                         process: SparkProcess,
-                                        createInteractiveRequest: CreateInteractiveRequest)
-  extends InteractiveWebSession(id, createInteractiveRequest) {
+                                        request: CreateInteractiveRequest)
+  extends InteractiveWebSession(id, process, request) {
 
   val stdoutThread = new Thread {
     override def run() = {
@@ -81,29 +75,4 @@ private class InteractiveSessionProcess(id: Int,
   stdoutThread.setName("process session stdout reader")
   stdoutThread.setDaemon(true)
   stdoutThread.start()
-
-  // Error out the job if the process errors out.
-  Future {
-    if (process.waitFor() != 0) {
-      _state = SessionState.Error()
-    } else {
-      // Set the state to done if the session shut down before contacting us.
-      _state match {
-        case (SessionState.Dead(_) | SessionState.Error(_) | SessionState.Success(_)) =>
-        case _ =>
-          _state = SessionState.Success()
-      }
-    }
-  }
-
-  override def logLines() = process.inputLines
-
-  override def stop(): Future[Unit] = {
-    super.stop().andThen { case r =>
-      // Make sure the process is reaped.
-      process.waitFor()
-      stdoutThread.join()
-      r
-    }
-  }
 }

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

@@ -23,12 +23,6 @@ import com.cloudera.hue.livy.spark.{SparkProcess, SparkProcessBuilderFactory}
 
 import scala.concurrent.ExecutionContext
 
-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"
-}
-
 class InteractiveSessionProcessFactory(processFactory: SparkProcessBuilderFactory)
   extends InteractiveSessionFactory(processFactory) {
 
@@ -42,15 +36,7 @@ class InteractiveSessionProcessFactory(processFactory: SparkProcessBuilderFactor
     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
   }
 }

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

@@ -45,14 +45,7 @@ private class InteractiveSessionYarn(id: Int,
                                      client: Client,
                                      process: SparkProcess,
                                      request: CreateInteractiveRequest)
-  extends InteractiveWebSession(id, request) {
-
-  // Error out the job if the process errors out.
-  Future {
-    if (process.waitFor() != 0) {
-      _state = SessionState.Error()
-    }
-  }
+  extends InteractiveWebSession(id, process, request) {
 
   private val job = Future {
     val job = client.getJobFromProcess(process)
@@ -67,8 +60,6 @@ private class InteractiveSessionYarn(id: Int,
   override def logLines() = process.inputLines
 
   override def stop(): Future[Unit] = {
-    process.destroy()
-
     super.stop().andThen {
       case _ =>
         try {

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

@@ -25,6 +25,7 @@ import com.cloudera.hue.livy._
 import com.cloudera.hue.livy.msgs.ExecuteRequest
 import com.cloudera.hue.livy.sessions._
 import com.cloudera.hue.livy.sessions.interactive.{Statement, InteractiveSession}
+import com.cloudera.hue.livy.spark.SparkProcess
 import dispatch._
 import org.json4s.JsonAST.{JNull, JString}
 import org.json4s.jackson.Serialization.write
@@ -34,7 +35,11 @@ import scala.annotation.tailrec
 import scala.concurrent.duration.Duration
 import scala.concurrent.{Future, _}
 
-abstract class InteractiveWebSession(val id: Int, request: CreateInteractiveRequest) extends InteractiveSession with Logging {
+abstract class InteractiveWebSession(val id: Int,
+                                     process: SparkProcess,
+                                     request: CreateInteractiveRequest)
+  extends InteractiveSession
+  with Logging {
 
   protected implicit def executor: ExecutionContextExecutor = ExecutionContext.global
   protected implicit def jsonFormats: Formats = DefaultFormats
@@ -49,6 +54,8 @@ abstract class InteractiveWebSession(val id: Int, request: CreateInteractiveRequ
 
   override def kind = request.kind
 
+  override def logLines() = process.inputLines
+
   override def proxyUser = request.proxyUser
 
   override def url: Option[URL] = _url
@@ -143,7 +150,7 @@ abstract class InteractiveWebSession(val id: Int, request: CreateInteractiveRequ
   }
 
   override def stop(): Future[Unit] = {
-    synchronized {
+    val future: Future[Unit] = synchronized {
       _state match {
         case SessionState.Idle() =>
           _state = SessionState.Busy()
@@ -185,6 +192,11 @@ abstract class InteractiveWebSession(val id: Int, request: CreateInteractiveRequ
           Future.successful(Unit)
       }
     }
+
+    future.andThen { case r =>
+      process.waitFor()
+      r
+    }
   }
 
   private def transition(state: SessionState) = synchronized {
@@ -215,4 +227,18 @@ abstract class InteractiveWebSession(val id: Int, request: CreateInteractiveRequ
       }
     }
   }
+
+  // Error out the job if the process errors out.
+  Future {
+    if (process.waitFor() == 0) {
+      // Set the state to done if the session shut down before contacting us.
+      _state match {
+        case (SessionState.Dead(_) | SessionState.Error(_) | SessionState.Success(_)) =>
+        case _ =>
+          _state = SessionState.Success()
+      }
+    } else {
+      _state = SessionState.Error()
+    }
+  }
 }