Browse Source

HUE-2865 [livy] Mark session as errored if session process died

Erick Tryzelaar 10 năm trước cách đây
mục cha
commit
243c7002c0

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

@@ -21,6 +21,7 @@ package com.cloudera.hue.livy.server.interactive
 import java.lang.ProcessBuilder.Redirect
 import java.net.URL
 
+import com.cloudera.hue.livy.sessions.Error
 import com.cloudera.hue.livy.spark.{SparkProcess, SparkSubmitProcessBuilder}
 import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.{RelativePath, AbsolutePath}
 import com.cloudera.hue.livy.{LivyConf, Logging, Utils}
@@ -119,14 +120,20 @@ private class InteractiveSessionProcess(id: Int,
   stdoutThread.setDaemon(true)
   stdoutThread.start()
 
+  // Error out the job if the process errors out.
+  Future {
+    if (process.waitFor() != 0) {
+      _state = Error()
+    }
+  }
+
   override def logLines() = process.inputLines
 
   override def stop(): Future[Unit] = {
-    super.stop() andThen { case r =>
+    super.stop().andThen { case r =>
       // Make sure the process is reaped.
       process.waitFor()
       stdoutThread.join()
-
       r
     }
   }

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

@@ -81,6 +81,13 @@ private class InteractiveSessionYarn(id: Int,
                                      createInteractiveRequest: CreateInteractiveRequest)
   extends InteractiveWebSession(id, createInteractiveRequest) {
 
+  // Error out the job if the process errors out.
+  Future {
+    if (process.waitFor() != 0) {
+      _state = Error()
+    }
+  }
+
   private val job = Future {
     val job = client.getJobFromProcess(process)