Browse Source

[livy] Move livy-yarn Job and ApplicationState into their own files

Erick Tryzelaar 10 years ago
parent
commit
a7f35e9290

+ 6 - 7
apps/spark/java/livy-server/src/main/scala/com/cloudera/hue/livy/server/batch/BatchSessionYarn.scala

@@ -19,16 +19,15 @@
 package com.cloudera.hue.livy.server.batch
 package com.cloudera.hue.livy.server.batch
 
 
 import java.lang.ProcessBuilder.Redirect
 import java.lang.ProcessBuilder.Redirect
+
 import com.cloudera.hue.livy.sessions._
 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.spark.SparkSubmitProcessBuilder
+import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.RelativePath
 import com.cloudera.hue.livy.yarn._
 import com.cloudera.hue.livy.yarn._
+import com.cloudera.hue.livy.{LineBufferedProcess, LivyConf}
 
 
-import scala.Error
 import scala.annotation.tailrec
 import scala.annotation.tailrec
-import scala.concurrent.{ExecutionContextExecutor, ExecutionContext, Future}
-import scala.util
+import scala.concurrent.{ExecutionContext, ExecutionContextExecutor, Future}
 
 
 object BatchSessionYarn {
 object BatchSessionYarn {
 
 
@@ -92,9 +91,9 @@ private class BatchSessionYarn(val id: Int, process: LineBufferedProcess, jobFut
             if (_state == SessionState.Running()) {
             if (_state == SessionState.Running()) {
               Thread.sleep(5000)
               Thread.sleep(5000)
               job.getStatus match {
               job.getStatus match {
-                case Client.SuccessfulFinish() =>
+                case ApplicationState.SuccessfulFinish() =>
                   _state = SessionState.Success(System.currentTimeMillis())
                   _state = SessionState.Success(System.currentTimeMillis())
-                case Client.UnsuccessfulFinish() =>
+                case ApplicationState.UnsuccessfulFinish() =>
                   _state = SessionState.Error(System.currentTimeMillis())
                   _state = SessionState.Error(System.currentTimeMillis())
                 case _ => aux()
                 case _ => aux()
               }
               }

+ 11 - 0
apps/spark/java/livy-yarn/src/main/scala/com/cloudera/hue/livy/yarn/ApplicationState.scala

@@ -0,0 +1,11 @@
+package com.cloudera.hue.livy.yarn
+
+sealed trait ApplicationState
+
+object ApplicationState {
+  case class New() extends ApplicationState
+  case class Accepted() extends ApplicationState
+  case class Running() extends ApplicationState
+  case class SuccessfulFinish() extends ApplicationState
+  case class UnsuccessfulFinish() extends ApplicationState
+}

+ 3 - 94
apps/spark/java/livy-yarn/src/main/scala/com/cloudera/hue/livy/yarn/Client.scala

@@ -18,29 +18,19 @@
 
 
 package com.cloudera.hue.livy.yarn
 package com.cloudera.hue.livy.yarn
 
 
-import java.io.{File, InputStream, BufferedReader, InputStreamReader}
+import java.io.File
 
 
-import com.cloudera.hue.livy.yarn.Client._
-import com.cloudera.hue.livy.{Utils, LineBufferedProcess, LivyConf, Logging}
-import org.apache.hadoop.yarn.api.records.{ApplicationId, FinalApplicationStatus, YarnApplicationState}
+import com.cloudera.hue.livy.{LineBufferedProcess, LivyConf, Logging}
+import org.apache.hadoop.fs.Path
 import org.apache.hadoop.yarn.client.api.YarnClient
 import org.apache.hadoop.yarn.client.api.YarnClient
 import org.apache.hadoop.yarn.conf.YarnConfiguration
 import org.apache.hadoop.yarn.conf.YarnConfiguration
 import org.apache.hadoop.yarn.util.ConverterUtils
 import org.apache.hadoop.yarn.util.ConverterUtils
-import org.apache.hadoop.fs.Path
 
 
 import scala.annotation.tailrec
 import scala.annotation.tailrec
 import scala.concurrent.ExecutionContext
 import scala.concurrent.ExecutionContext
-import scala.io.Source
 
 
 object Client {
 object Client {
   private lazy val regex = """Application report for (\w+)""".r.unanchored
   private lazy val regex = """Application report for (\w+)""".r.unanchored
-
-  sealed trait ApplicationStatus
-  case class New() extends ApplicationStatus
-  case class Accepted() extends ApplicationStatus
-  case class Running() extends ApplicationStatus
-  case class SuccessfulFinish() extends ApplicationStatus
-  case class UnsuccessfulFinish() extends ApplicationStatus
 }
 }
 
 
 class FailedToSubmitApplication extends Exception
 class FailedToSubmitApplication extends Exception
@@ -85,85 +75,4 @@ class Client(livyConf: LivyConf) extends Logging {
   }
   }
 }
 }
 
 
-class Job(yarnClient: YarnClient, appId: ApplicationId) {
-  def waitForFinish(timeoutMs: Long): Option[ApplicationStatus] = {
-    val startTimeMs = System.currentTimeMillis()
-
-    while (System.currentTimeMillis() - startTimeMs < timeoutMs) {
-      val status = getStatus
-      status match {
-        case SuccessfulFinish() | UnsuccessfulFinish() =>
-          return Some(status)
-        case _ =>
-      }
-
-      Thread.sleep(1000)
-    }
-
-    None
-  }
-
-  def waitForStatus(status: ApplicationStatus, timeoutMs: Long): Option[ApplicationStatus] = {
-    val startTimeMs = System.currentTimeMillis()
-
-    while (System.currentTimeMillis() - startTimeMs < timeoutMs) {
-      if (getStatus == status) {
-        return Some(status)
-      }
-
-      Thread.sleep(1000)
-    }
-
-    None
-  }
-
-  def waitForRPC(timeoutMs: Long): Option[(String, Int)] = {
-    waitForStatus(Running(), timeoutMs)
-
-    val startTimeMs = System.currentTimeMillis()
 
 
-    while (System.currentTimeMillis() - startTimeMs < timeoutMs) {
-      val statusResponse = yarnClient.getApplicationReport(appId)
-
-      (statusResponse.getHost, statusResponse.getRpcPort) match {
-        case ("N/A", _) | (_, -1) =>
-        case (hostname, port) => return Some((hostname, port))
-      }
-    }
-
-    None
-  }
-
-  def getHost: String = {
-    val statusResponse = yarnClient.getApplicationReport(appId)
-    statusResponse.getHost
-  }
-
-  def getPort: Int = {
-    val statusResponse = yarnClient.getApplicationReport(appId)
-    statusResponse.getRpcPort
-  }
-
-  def getStatus: ApplicationStatus = {
-    val statusResponse = yarnClient.getApplicationReport(appId)
-    convertState(statusResponse.getYarnApplicationState, statusResponse.getFinalApplicationStatus)
-  }
-
-  def stop(): Unit = {
-    yarnClient.killApplication(appId)
-  }
-
-  private def convertState(state: YarnApplicationState, status: FinalApplicationStatus): ApplicationStatus = {
-    (state, status) match {
-      case (YarnApplicationState.FINISHED, FinalApplicationStatus.SUCCEEDED) => SuccessfulFinish()
-      case (YarnApplicationState.FINISHED, _) |
-           (YarnApplicationState.KILLED, _) |
-           (YarnApplicationState.FAILED, _) => UnsuccessfulFinish()
-      case (YarnApplicationState.NEW, _) |
-           (YarnApplicationState.NEW_SAVING, _) |
-           (YarnApplicationState.SUBMITTED, _) => New()
-      case (YarnApplicationState.RUNNING, _) => Running()
-      case (YarnApplicationState.ACCEPTED, _) => Accepted()
-    }
-  }
-}

+ 87 - 0
apps/spark/java/livy-yarn/src/main/scala/com/cloudera/hue/livy/yarn/Job.scala

@@ -0,0 +1,87 @@
+package com.cloudera.hue.livy.yarn
+
+import org.apache.hadoop.yarn.api.records.{FinalApplicationStatus, YarnApplicationState, ApplicationId}
+import org.apache.hadoop.yarn.client.api.YarnClient
+
+class Job(yarnClient: YarnClient, appId: ApplicationId) {
+  def waitForFinish(timeoutMs: Long): Option[ApplicationState] = {
+    val startTimeMs = System.currentTimeMillis()
+
+    while (System.currentTimeMillis() - startTimeMs < timeoutMs) {
+      val status = getStatus
+      status match {
+        case ApplicationState.SuccessfulFinish() | ApplicationState.UnsuccessfulFinish() =>
+          return Some(status)
+        case _ =>
+      }
+
+      Thread.sleep(1000)
+    }
+
+    None
+  }
+
+  def waitForStatus(status: ApplicationState, timeoutMs: Long): Option[ApplicationState] = {
+    val startTimeMs = System.currentTimeMillis()
+
+    while (System.currentTimeMillis() - startTimeMs < timeoutMs) {
+      if (getStatus == status) {
+        return Some(status)
+      }
+
+      Thread.sleep(1000)
+    }
+
+    None
+  }
+
+  def waitForRPC(timeoutMs: Long): Option[(String, Int)] = {
+    waitForStatus(ApplicationState.Running(), timeoutMs)
+
+    val startTimeMs = System.currentTimeMillis()
+
+    while (System.currentTimeMillis() - startTimeMs < timeoutMs) {
+      val statusResponse = yarnClient.getApplicationReport(appId)
+
+      (statusResponse.getHost, statusResponse.getRpcPort) match {
+        case ("N/A", _) | (_, -1) =>
+        case (hostname, port) => return Some((hostname, port))
+      }
+    }
+
+    None
+  }
+
+  def getHost: String = {
+    val statusResponse = yarnClient.getApplicationReport(appId)
+    statusResponse.getHost
+  }
+
+  def getPort: Int = {
+    val statusResponse = yarnClient.getApplicationReport(appId)
+    statusResponse.getRpcPort
+  }
+
+  def getStatus: ApplicationState = {
+    val statusResponse = yarnClient.getApplicationReport(appId)
+    convertState(statusResponse.getYarnApplicationState, statusResponse.getFinalApplicationStatus)
+  }
+
+  def stop(): Unit = {
+    yarnClient.killApplication(appId)
+  }
+
+  private def convertState(state: YarnApplicationState, status: FinalApplicationStatus): ApplicationState = {
+    (state, status) match {
+      case (YarnApplicationState.FINISHED, FinalApplicationStatus.SUCCEEDED) => ApplicationState.SuccessfulFinish()
+      case (YarnApplicationState.FINISHED, _) |
+           (YarnApplicationState.KILLED, _) |
+           (YarnApplicationState.FAILED, _) => ApplicationState.UnsuccessfulFinish()
+      case (YarnApplicationState.NEW, _) |
+           (YarnApplicationState.NEW_SAVING, _) |
+           (YarnApplicationState.SUBMITTED, _) => ApplicationState.New()
+      case (YarnApplicationState.RUNNING, _) => ApplicationState.Running()
+      case (YarnApplicationState.ACCEPTED, _) => ApplicationState.Accepted()
+    }
+  }
+}