Эх сурвалжийг харах

[spark] Cleanup livy-server

Erick Tryzelaar 11 жил өмнө
parent
commit
5831d8560a

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

@@ -2,15 +2,10 @@ package com.cloudera.hue.livy.server
 
 import javax.servlet.ServletContext
 
-import scala.concurrent.duration._
 import com.cloudera.hue.livy.WebServer
-import org.json4s.{DefaultFormats, Formats}
 import org.scalatra._
-import org.scalatra.json.JacksonJsonSupport
 import org.scalatra.servlet.ScalatraListener
 
-import scala.concurrent.{Await, ExecutionContext, ExecutionContextExecutor}
-
 object Main {
   def main(args: Array[String]): Unit = {
     val port = sys.env.getOrElse("PORT", "8998").toInt
@@ -39,79 +34,3 @@ class ScalatraBootstrap extends LifeCycle {
     sessionManager.close()
   }
 }
-
-class WebApp(sessionManager: SessionManager) extends ScalatraServlet with FutureSupport with MethodOverride with JacksonJsonSupport with UrlGeneratorSupport {
-
-  override protected implicit def executor: ExecutionContextExecutor = ExecutionContext.global
-  override protected implicit def jsonFormats: Formats = DefaultFormats
-
-  before() {
-    contentType = formats("json")
-  }
-
-  get("/sessions") {
-    sessionManager.getSessionIds
-  }
-
-  post("/sessions") {
-    val createSessionRequest = parsedBody.extract[CreateSessionRequest]
-
-    val sessionFuture = createSessionRequest.lang match {
-      case "scala" => sessionManager.createSparkSession()
-      case lang => halt(400, "unsupported language: " + lang)
-    }
-
-    val rep = for {
-      session <- sessionFuture
-    } yield redirect(url(getSession, "sessionId" -> session.id))
-
-    new AsyncResult { val is = rep }
-  }
-
-  val getStatements = get("/sessions/:sessionId/statements") {
-    sessionManager.get(params("sessionId")) match {
-      case Some(session: Session) =>
-        val statements = session.statements()
-        val statementsWaited = Await.result(statements, Duration.Inf) //5 seconds)
-        //new AsyncResult() { val is = statements }
-        statementsWaited
-      case None => NotFound("Session not found")
-    }
-  }
-
-  val getSession = get("/sessions/:sessionId") {
-    redirect(url(getStatements, "sessionId" -> params("sessionId")))
-  }
-
-  delete("/sessions/:sessionId") {
-    sessionManager.close(params("sessionId"))
-    NoContent
-  }
-
-  post("/sessions/:sessionId/statements") {
-    val req = parsedBody.extract[ExecuteStatementRequest]
-
-    sessionManager.get(params("sessionId")) match {
-      case Some(session) =>
-        val statement = session.executeStatement(req.statement)
-        val foo = Await.result(statement, Duration.Inf)
-        foo
-
-
-        //new AsyncResult() { val is = statement }
-      case None => NotFound("Session not found")
-    }
-  }
-
-  val getStatement = get("/sessions/:sessionId/statements/:statementId") {
-    sessionManager.get(params("sessionId")) match {
-      case Some(session) =>
-        val statement = session.statement(params("statementId").toInt)
-        new AsyncResult() { val is = statement }
-      case None => NotFound("Session not found")
-    }
-  }
-}
-
-case class CreateSessionRequest(lang: String)
-case class ExecuteStatementRequest(statement: String)

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

@@ -14,7 +14,7 @@ class ProcessSessionFactory extends SessionFactory {
 
   override def createSparkSession: Future[Session] = {
     future {
-      val id = "a" //UUID.randomUUID().toString
+      val id = UUID.randomUUID().toString
       new SparkProcessSession(id)
     }
   }

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

@@ -0,0 +1,86 @@
+package com.cloudera.hue.livy.server
+
+import org.json4s.{DefaultFormats, Formats}
+import org.scalatra._
+import org.scalatra.json.JacksonJsonSupport
+
+import scala.concurrent.{ExecutionContext, ExecutionContextExecutor}
+
+object WebApp {
+  case class CreateSessionRequest(lang: String)
+  case class ExecuteStatementRequest(statement: String)
+}
+
+class WebApp(sessionManager: SessionManager)
+  extends ScalatraServlet
+  with FutureSupport
+  with MethodOverride
+  with JacksonJsonSupport
+  with UrlGeneratorSupport {
+
+  import com.cloudera.hue.livy.server.WebApp._
+
+  override protected implicit def executor: ExecutionContextExecutor = ExecutionContext.global
+  override protected implicit def jsonFormats: Formats = DefaultFormats
+
+  before() {
+    contentType = formats("json")
+  }
+
+  get("/sessions") {
+    sessionManager.getSessionIds
+  }
+
+  post("/sessions") {
+    val createSessionRequest = parsedBody.extract[CreateSessionRequest]
+
+    val sessionFuture = createSessionRequest.lang match {
+      case "scala" => sessionManager.createSparkSession()
+      case lang => halt(400, "unsupported language: " + lang)
+    }
+
+    val rep = for {
+      session <- sessionFuture
+    } yield redirect(url(getSession, "sessionId" -> session.id))
+
+    new AsyncResult { val is = rep }
+  }
+
+  val getStatements = get("/sessions/:sessionId/statements") {
+    sessionManager.get(params("sessionId")) match {
+      case Some(session: Session) =>
+        val statements = session.statements()
+        new AsyncResult() { val is = statements }
+      case None => NotFound("Session not found")
+    }
+  }
+
+  val getSession = get("/sessions/:sessionId") {
+    redirect(url(getStatements, "sessionId" -> params("sessionId")))
+  }
+
+  delete("/sessions/:sessionId") {
+    sessionManager.close(params("sessionId"))
+    NoContent
+  }
+
+  post("/sessions/:sessionId/statements") {
+    val req = parsedBody.extract[ExecuteStatementRequest]
+
+    sessionManager.get(params("sessionId")) match {
+      case Some(session) =>
+        val statement = session.executeStatement(req.statement)
+        new AsyncResult() { val is = statement }
+      case None => NotFound("Session not found")
+    }
+  }
+
+  val getStatement = get("/sessions/:sessionId/statements/:statementId") {
+    sessionManager.get(params("sessionId")) match {
+      case Some(session) =>
+        val statement = session.statement(params("statementId").toInt)
+        new AsyncResult() { val is = statement }
+      case None => NotFound("Session not found")
+    }
+  }
+}