浏览代码

[livy] Cut down on the number of Futures

This makes the code a little easier to debug.
Erick Tryzelaar 10 年之前
父节点
当前提交
48db6e3

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

@@ -26,7 +26,7 @@ abstract class SessionFactory[S <: Session] {
 
   protected implicit def jsonFormats: Formats = DefaultFormats
 
-  def create(id: Int, createRequest: JValue): Future[S]
+  def create(id: Int, createRequest: JValue): S
 
   def close(): Unit = {}
 }

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

@@ -47,15 +47,16 @@ class SessionManager[S <: Session](factory: SessionFactory[S])
   private val garbageCollector = new GarbageCollector
   garbageCollector.start()
 
-  def create(createRequest: JValue): Future[S] = synchronized {
+  def create(createRequest: JValue): S = {
     val id = _idCounter.getAndIncrement
-    val session: Future[S] = factory.create(id, createRequest)
+    val session: S = factory.create(id, createRequest)
 
-    session.map({ case (session) =>
-      info("created session %s" format session.id)
+    info("created session %s" format session.id)
+
+    synchronized {
       _sessions.put(session.id, session)
       session
-    })
+    }
   }
 
   def get(id: Int): Option[S] = _sessions.get(id)

+ 5 - 4
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/SessionServlet.scala

@@ -26,7 +26,7 @@ import org.json4s.{MappingException, DefaultFormats, Formats, JValue}
 import org.scalatra._
 import org.scalatra.json.JacksonJsonSupport
 
-import scala.concurrent.ExecutionContext
+import scala.concurrent.{Future, ExecutionContext}
 
 object SessionServlet extends Logging
 
@@ -109,11 +109,12 @@ abstract class SessionServlet[S <: Session](sessionManager: SessionManager[S])
 
   post("/") {
     new AsyncResult {
-      val is = for {
-        session <- sessionManager.create(parsedBody)
-      } yield Created(session,
+      val is = Future {
+        val session = sessionManager.create(parsedBody)
+        Created(session,
           headers = Map("Location" -> url(getSession, "id" -> session.id.toString))
         )
+      }
     }
   }
 

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

@@ -27,5 +27,5 @@ abstract class BatchSessionFactory extends SessionFactory[BatchSession] {
   override def create(id: Int, createRequest: JValue) =
     create(id, createRequest.extract[CreateBatchRequest])
 
-  def create(id: Int, createRequest: CreateBatchRequest): Future[BatchSession]
+  def create(id: Int, createRequest: CreateBatchRequest): BatchSession
 }

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

@@ -25,6 +25,6 @@ import scala.concurrent.Future
 class BatchSessionProcessFactory(livyConf: LivyConf)
   extends BatchSessionFactory
 {
-  override def create(id: Int, createBatchRequest: CreateBatchRequest): Future[BatchSession] =
-    Future.successful(BatchSessionProcess(livyConf, id, createBatchRequest))
+  override def create(id: Int, createBatchRequest: CreateBatchRequest): BatchSession =
+    BatchSessionProcess(livyConf, id, createBatchRequest)
 }

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

@@ -27,6 +27,6 @@ class BatchSessionYarnFactory(livyConf: LivyConf) extends BatchSessionFactory {
 
   val client = new Client(livyConf)
 
-  def create(id: Int, createBatchRequest: CreateBatchRequest): Future[BatchSession] =
-    Future.successful(BatchSessionYarn(livyConf, client, id, createBatchRequest))
+  def create(id: Int, createBatchRequest: CreateBatchRequest): BatchSession =
+    BatchSessionYarn(livyConf, client, id, createBatchRequest)
 }

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

@@ -31,5 +31,5 @@ trait InteractiveSessionFactory extends SessionFactory[InteractiveSession] {
   override def create(id: Int, createRequest: JValue) =
     create(id, createRequest.extract[CreateInteractiveRequest])
 
-  def create(id: Int, createRequest: CreateInteractiveRequest): Future[InteractiveSession]
+  def create(id: Int, createRequest: CreateInteractiveRequest): InteractiveSession
 }

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

@@ -26,9 +26,7 @@ class InteractiveSessionProcessFactory(livyConf: LivyConf) extends InteractiveSe
 
    implicit def executor: ExecutionContext = ExecutionContext.global
 
-   override def create(id: Int, createInteractiveRequest: CreateInteractiveRequest): Future[InteractiveSession] = {
-     Future {
-       InteractiveSessionProcess.create(livyConf, id, createInteractiveRequest)
-     }
+   override def create(id: Int, createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
+     InteractiveSessionProcess.create(livyConf, id, createInteractiveRequest)
    }
  }

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

@@ -29,10 +29,8 @@ class InteractiveSessionYarnFactory(livyConf: LivyConf) extends InteractiveSessi
 
    val client = new Client(livyConf)
 
-   override def create(id: Int, createInteractiveRequest: CreateInteractiveRequest): Future[InteractiveSession] = {
-     Future {
-       InteractiveSessionYarn.create(livyConf, client, id, createInteractiveRequest)
-     }
+   override def create(id: Int, createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
+     InteractiveSessionYarn.create(livyConf, client, id, createInteractiveRequest)
    }
 
    override def close(): Unit = {

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

@@ -75,8 +75,8 @@ class InteractiveSessionServletSpec extends ScalatraSuite with FunSpecLike {
   }
 
   class MockInteractiveSessionFactory() extends InteractiveSessionFactory {
-    override def create(id: Int, createInteractiveRequest: CreateInteractiveRequest): Future[InteractiveSession] = {
-      Future.successful(new MockInteractiveSession(id))
+    override def create(id: Int, createInteractiveRequest: CreateInteractiveRequest): InteractiveSession = {
+      new MockInteractiveSession(id)
     }
   }