فهرست منبع

[livy] Start merging batch and interactive sessions

Erick Tryzelaar 10 سال پیش
والد
کامیت
af17ebd

+ 8 - 0
apps/spark/java/livy-core/src/main/scala/com/cloudera/hue/livy/sessions/State.scala

@@ -32,6 +32,10 @@ case class Idle() extends State {
   override def toString = "idle"
 }
 
+case class Running() extends State {
+  override def toString = "running"
+}
+
 case class Busy() extends State {
   override def toString = "busy"
 }
@@ -47,3 +51,7 @@ case class ShuttingDown() extends State {
 case class Dead() extends State {
   override def toString = "dead"
 }
+
+case class Success() extends State {
+  override def toString = "success"
+}

+ 2 - 0
apps/spark/java/livy-repl/src/main/scala/com/cloudera/hue/livy/repl/WebApp.scala

@@ -50,9 +50,11 @@ class WebApp(session: Session) extends ScalatraServlet with FutureSupport with J
       case Starting() => "starting"
       case Idle() => "idle"
       case Busy() => "busy"
+      case Running() => "running"
       case Error() => "error"
       case ShuttingDown() => "shutting_down"
       case Dead() => "dead"
+      case Success() => "success"
     }
     Map("state" -> state)
   }

+ 5 - 15
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/batch/State.scala → apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/Session.scala

@@ -16,22 +16,12 @@
  * limitations under the License.
  */
 
-package com.cloudera.hue.livy.server.batch
+package com.cloudera.hue.livy.server
 
-sealed trait State
+import com.cloudera.hue.livy.sessions.State
 
-case class Starting() extends State {
-  override def toString = "starting"
-}
-
-case class Running() extends State {
-  override def toString = "running"
-}
-
-case class Error() extends State {
-  override def toString = "error"
-}
+trait Session {
+  def id: Int
 
-case class Success() extends State {
-  override def toString = "success"
+  def state: State
 }

+ 3 - 5
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/batch/Batch.scala

@@ -18,13 +18,11 @@
 
 package com.cloudera.hue.livy.server.batch
 
-import scala.concurrent.Future
-
-abstract class Batch {
-  def id: Int
+import com.cloudera.hue.livy.server.Session
 
-  def state: State
+import scala.concurrent.Future
 
+trait Batch extends Session {
   def lines: IndexedSeq[String]
 
   def stop(): Future[Unit]

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

@@ -20,6 +20,7 @@ package com.cloudera.hue.livy.server.batch
 
 import java.lang.ProcessBuilder.Redirect
 
+import com.cloudera.hue.livy.sessions.{Success, Running, State}
 import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.RelativePath
 import com.cloudera.hue.livy.{LivyConf, LineBufferedProcess}
 import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder

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

@@ -20,11 +20,13 @@ package com.cloudera.hue.livy.server.batch
 
 import java.lang.ProcessBuilder.Redirect
 
+import com.cloudera.hue.livy.sessions._
 import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.RelativePath
 import com.cloudera.hue.livy.{LineBufferedProcess, LivyConf}
 import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder
 import com.cloudera.hue.livy.yarn._
 
+import scala.Error
 import scala.annotation.tailrec
 import scala.concurrent.{ExecutionContextExecutor, ExecutionContext, Future}
 import scala.util

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

@@ -23,6 +23,7 @@ import java.util.concurrent.TimeoutException
 
 import com.cloudera.hue.livy.Utils
 import com.cloudera.hue.livy.msgs.ExecuteRequest
+import com.cloudera.hue.livy.server.Session
 import com.cloudera.hue.livy.sessions.{Kind, State}
 
 import scala.concurrent._
@@ -34,17 +35,13 @@ object InteractiveSession {
   class StatementNotFound extends Exception
 }
 
-trait InteractiveSession {
-  def id: Int
-
+trait InteractiveSession extends Session {
   def kind: Kind
 
   def proxyUser: Option[String]
 
   def lastActivity: Long
 
-  def state: State
-
   def url: Option[URL]
 
   def url_=(url: URL)

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

@@ -118,7 +118,7 @@ class InteractiveWebSession(val id: Int,
             waitForStateChange(Starting(), Duration(10, TimeUnit.SECONDS))
             stop()
           }
-        case Busy() =>
+        case Busy() | Running() =>
           Future {
             waitForStateChange(Busy(), Duration(10, TimeUnit.SECONDS))
             stop()
@@ -128,7 +128,7 @@ class InteractiveWebSession(val id: Int,
             waitForStateChange(ShuttingDown(), Duration(10, TimeUnit.SECONDS))
             stop()
           }
-        case Error() | Dead() =>
+        case Error() | Dead() | Success() =>
           Future.successful(Unit)
       }
     }

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

@@ -22,8 +22,9 @@ import java.io.FileWriter
 import java.nio.file.{Files, Path}
 import java.util.concurrent.TimeUnit
 
+import com.cloudera.hue.livy.sessions.Success
 import com.cloudera.hue.livy.{LivyConf, Utils}
-import com.cloudera.hue.livy.server.batch.{Success, CreateBatchRequest, BatchProcess}
+import com.cloudera.hue.livy.server.batch.{CreateBatchRequest, BatchProcess}
 import org.scalatest.{ShouldMatchers, BeforeAndAfterAll, FunSpec}
 
 import scala.concurrent.duration.Duration

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

@@ -22,6 +22,7 @@ import java.io.FileWriter
 import java.nio.file.{Files, Path}
 import java.util.concurrent.TimeUnit
 
+import com.cloudera.hue.livy.sessions.Success
 import com.cloudera.hue.livy.{LivyConf, Utils}
 import com.cloudera.hue.livy.server.batch._
 import org.json4s.JsonAST.{JArray, JInt, JObject, JString}