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

[spark] Add livy-core, cleanup launching livy-repl

Erick Tryzelaar 11 жил өмнө
parent
commit
80380770ae

+ 49 - 0
apps/spark/java/livy-core/pom.xml

@@ -0,0 +1,49 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<project xmlns="http://maven.apache.org/POM/4.0.0"
+         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
+         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
+    <modelVersion>4.0.0</modelVersion>
+    <parent>
+        <groupId>com.cloudera.hue.livy</groupId>
+        <artifactId>livy-main</artifactId>
+        <relativePath>../pom.xml</relativePath>
+        <version>3.7.0-SNAPSHOT</version>
+    </parent>
+
+    <artifactId>livy-core</artifactId>
+    <packaging>jar</packaging>
+
+    <build>
+        <plugins>
+
+            <plugin>
+                <groupId>org.scala-tools</groupId>
+                <artifactId>maven-scala-plugin</artifactId>
+                <version>2.15.2</version>
+                <executions>
+                    <execution>
+                        <goals>
+                            <goal>compile</goal>
+                            <goal>testCompile</goal>
+                        </goals>
+                    </execution>
+                </executions>
+            </plugin>
+
+        </plugins>
+    </build>
+
+    <reporting>
+        <plugins>
+            <plugin>
+                <groupId>org.scala-tools</groupId>
+                <artifactId>maven-scala-plugin</artifactId>
+                <configuration>
+                    <scalaVersion>${scala.version}</scalaVersion>
+                </configuration>
+            </plugin>
+        </plugins>
+    </reporting>
+
+</project>
+

+ 1 - 3
apps/spark/java/livy-yarn/src/main/scala/com/cloudera/hue/livy/yarn/Logging.scala → apps/spark/java/livy-core/src/main/scala/com/cloudera/hue/livy/Logging.scala

@@ -1,6 +1,4 @@
-package com.cloudera.hue.livy.yarn
-
-import org.slf4j.LoggerFactory
+package com.cloudera.hue.livy
 
 trait Logging {
   val loggerName = this.getClass.getName

+ 6 - 0
apps/spark/java/livy-repl/pom.xml

@@ -66,6 +66,12 @@
             <scope>provided</scope>
         </dependency>
 
+        <dependency>
+            <groupId>com.cloudera.hue.livy</groupId>
+            <artifactId>livy-core</artifactId>
+            <version>${project.version}</version>
+        </dependency>
+
     </dependencies>
 
     <build>

+ 0 - 20
apps/spark/java/livy-repl/src/main/scala/Scalatra.scala

@@ -1,20 +0,0 @@
-import javax.servlet.ServletContext
-
-import com.cloudera.hue.livy.repl.LivyApp
-import com.cloudera.hue.livy.repl.interpreter.SparkInterpreter
-import org.scalatra.LifeCycle
-
-class ScalatraBootstrap2 extends LifeCycle {
-
-  //val system = ActorSystem()
-  val sparkInterpreter = new SparkInterpreter
-
-  override def init(context: ServletContext): Unit = {
-    context.mount(new LivyApp(sparkInterpreter), "/*")
-  }
-
-  override def destroy(context: ServletContext): Unit = {
-    sparkInterpreter.close()
-    //system.shutdown()
-  }
-}

+ 4 - 19
apps/spark/java/livy-repl/src/main/scala/com/cloudera/hue/livy/repl/Main.scala

@@ -1,29 +1,14 @@
 package com.cloudera.hue.livy.repl
 
-import org.eclipse.jetty.server.Server
-import org.eclipse.jetty.servlet.{ServletHolder, DefaultServlet}
-import org.eclipse.jetty.webapp.WebAppContext
-import org.scalatra.servlet.{AsyncSupport, ScalatraListener}
-
-import scala.concurrent.ExecutionContext
-
 object Main {
   def main(args: Array[String]): Unit = {
     val port = sys.env.getOrElse("PORT", "8999").toInt
-    val server = new Server(port)
-    val context = new WebAppContext()
-
-    context.setContextPath("/")
-    context.setResourceBase("src/main/com/cloudera/hue/livy/repl")
-    context.addEventListener(new ScalatraListener)
-
-    context.addServlet(classOf[DefaultServlet], "/")
-
-    context.setAttribute(AsyncSupport.ExecutionContextKey, ExecutionContext.global)
-
-    server.setHandler(context)
+    val server = new WebServer(port)
+    server.start()
+    server.join()
 
     server.start()
     server.join()
+    server.stop()
   }
 }

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

@@ -0,0 +1,70 @@
+package com.cloudera.hue.livy.repl
+
+import javax.servlet.ServletContext
+
+import com.cloudera.hue.livy.Logging
+import com.cloudera.hue.livy.repl.interpreter.SparkInterpreter
+import org.eclipse.jetty.server.Server
+import org.eclipse.jetty.servlet.DefaultServlet
+import org.eclipse.jetty.webapp.WebAppContext
+import org.scalatra.servlet.{AsyncSupport, ScalatraListener}
+import org.scalatra.{LifeCycle, ScalatraServlet}
+
+import scala.concurrent.ExecutionContext
+
+class WebServer(var port: Int) extends Logging {
+  val server = new Server(port)
+  val context = new WebAppContext()
+
+  context.setContextPath("/")
+  context.setResourceBase("src/main/com/cloudera/hue/livy/repl")
+  context.setInitParameter(ScalatraListener.LifeCycleKey, classOf[ScalatraBootstrap].getCanonicalName)
+  context.addEventListener(new ScalatraListener)
+
+  context.addServlet(classOf[DefaultServlet], "/")
+
+  context.setAttribute(AsyncSupport.ExecutionContextKey, ExecutionContext.global)
+
+  server.setHandler(context)
+
+  def start() = {
+    server.start()
+    port = server.getConnectors()(0).getLocalPort
+
+    info("Starting RPC server on %s" format port)
+  }
+
+  def join() = {
+    server.join()
+  }
+
+  def stop() = {
+    context.stop()
+    server.stop()
+  }
+}
+
+class ScalatraBootstrap extends LifeCycle {
+
+  //val system = ActorSystem()
+  val sparkInterpreter = new SparkInterpreter
+
+  override def init(context: ServletContext): Unit = {
+    context.mount(new LivyApp(sparkInterpreter), "/*")
+  }
+
+  override def destroy(context: ServletContext): Unit = {
+    sparkInterpreter.close()
+    //system.shutdown()
+  }
+}
+
+class WebApp extends ScalatraServlet {
+  get("/") {
+    "hello world"
+  }
+
+  get("/hello") {
+    "hello world2"
+  }
+}

+ 6 - 19
apps/spark/java/livy-yarn/pom.xml

@@ -72,19 +72,6 @@
             <version>${scala.version}</version>
         </dependency>
 
-        <!--
-        <dependency>
-            <groupId>org.scalatra</groupId>
-            <artifactId>scalatra_2.10</artifactId>
-            <version>2.3.0</version>
-            <exclusions>
-                <exclusion>
-                    <groupId>org.scala-lang</groupId>
-                    <artifactId>scala-compiler</artifactId>
-                </exclusion>
-            </exclusions>
-        </dependency>
-        -->
 
         <dependency>
             <groupId>org.eclipse.jetty</groupId>
@@ -92,17 +79,17 @@
             <version>8.1.14.v20131031</version>
         </dependency>
 
+        <dependency>
+            <groupId>com.cloudera.hue.livy</groupId>
+            <artifactId>livy-core</artifactId>
+            <version>${project.version}</version>
+        </dependency>
+
         <dependency>
             <groupId>com.cloudera.hue.livy</groupId>
             <artifactId>livy-repl</artifactId>
             <version>${project.version}</version>
             <exclusions>
-                <!--
-                <exclusion>
-                    <groupId>org.scala-lang</groupId>
-                    <artifactId>scala-compiler</artifactId>
-                </exclusion>
-                -->
                 <exclusion>
                     <groupId>org.slf4j</groupId>
                     <artifactId>slf4j-api</artifactId>

+ 0 - 20
apps/spark/java/livy-yarn/src/main/scala/Scalatra.scala

@@ -1,20 +0,0 @@
-import javax.servlet.ServletContext
-
-import com.cloudera.hue.livy.repl.interpreter.SparkInterpreter
-import com.cloudera.hue.livy.repl.LivyApp
-import org.scalatra.LifeCycle
-
-class ScalatraBootstrap extends LifeCycle {
-
-  //val system = ActorSystem()
-  val sparkInterpreter = new SparkInterpreter
-
-  override def init(context: ServletContext): Unit = {
-    context.mount(new LivyApp(sparkInterpreter), "/*")
-  }
-
-  override def destroy(context: ServletContext): Unit = {
-    sparkInterpreter.close()
-    //system.shutdown()
-  }
-}

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

@@ -1,5 +1,7 @@
 package com.cloudera.hue.livy.yarn
 
+import com.cloudera.hue.livy.Logging
+import com.cloudera.hue.livy.repl.WebServer
 import org.apache.hadoop.yarn.api.ApplicationConstants
 import org.apache.hadoop.yarn.api.records.FinalApplicationStatus
 import org.apache.hadoop.yarn.client.api.AMRMClient
@@ -30,7 +32,7 @@ object AppMaster extends Logging {
 }
 
 class AppMasterService(yarnConfig: YarnConfiguration, nodeHostString: String) extends Logging {
-  val webServer = new WebServer
+  val webServer = new WebServer(0)
   val amRMClient = AMRMClient.createAMRMClient()
   amRMClient.init(yarnConfig)
 

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

@@ -1,5 +1,6 @@
 package com.cloudera.hue.livy.yarn
 
+import com.cloudera.hue.livy.Logging
 import org.apache.hadoop.fs.{FileSystem, Path}
 import org.apache.hadoop.yarn.api.ApplicationConstants
 import org.apache.hadoop.yarn.api.records._

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

@@ -1,49 +0,0 @@
-package com.cloudera.hue.livy.yarn
-
-import org.eclipse.jetty.server.Server
-import org.eclipse.jetty.servlet.{DefaultServlet, ServletHolder}
-import org.eclipse.jetty.webapp.WebAppContext
-import org.scalatra.ScalatraServlet
-import org.scalatra.servlet.{ScalatraListener, AsyncSupport}
-
-import scala.concurrent.ExecutionContext
-
-class WebServer extends Logging {
-  val server = new Server(0)
-  val context = new WebAppContext()
-  var port = 0
-
-  context.setContextPath("/")
-  context.setResourceBase("src/main/com/cloudera/hue/livy/yarn")
-  context.addEventListener(new ScalatraListener)
-
-  context.addServlet(classOf[DefaultServlet], "/")
-
-  context.setAttribute(AsyncSupport.ExecutionContextKey, ExecutionContext.global)
-
-  server.setHandler(context)
-
-  def start() = {
-    //context.setContextPath("/")
-    //context.setResourceBase(getClass.getClassLoader.getResource())
-    server.start()
-    port = server.getConnectors()(0).getLocalPort
-
-    info("Starting RPC server on %s" format port)
-  }
-
-  def stop() = {
-    context.stop()
-    server.stop()
-  }
-}
-
-class WebApp extends ScalatraServlet {
-  get("/") {
-    "hello world"
-  }
-
-  get("/hello") {
-    "hello world2"
-  }
-}

+ 1 - 0
apps/spark/java/pom.xml

@@ -56,6 +56,7 @@
         <module>livy-assembly</module>
         <module>livy-server</module>
         -->
+        <module>livy-core</module>
         <module>livy-repl</module>
         <module>livy-yarn</module>
     </modules>