|
@@ -2,6 +2,7 @@ package com.cloudera.hue.livy.repl.scala.interpreter
|
|
|
|
|
|
|
|
import java.io._
|
|
import java.io._
|
|
|
|
|
|
|
|
|
|
+import org.apache.spark.{SparkConf, SparkContext}
|
|
|
import org.apache.spark.repl.SparkIMain
|
|
import org.apache.spark.repl.SparkIMain
|
|
|
|
|
|
|
|
import scala.concurrent.ExecutionContext
|
|
import scala.concurrent.ExecutionContext
|
|
@@ -29,6 +30,7 @@ class Interpreter {
|
|
|
private var _state: Interpreter.State = Interpreter.NotStarted()
|
|
private var _state: Interpreter.State = Interpreter.NotStarted()
|
|
|
private val outputStream = new ByteArrayOutputStream()
|
|
private val outputStream = new ByteArrayOutputStream()
|
|
|
private var sparkIMain: SparkIMain = _
|
|
private var sparkIMain: SparkIMain = _
|
|
|
|
|
+ private var sparkContext: SparkContext = _
|
|
|
private var executeCount = 0
|
|
private var executeCount = 0
|
|
|
|
|
|
|
|
def state = _state
|
|
def state = _state
|
|
@@ -44,11 +46,24 @@ class Interpreter {
|
|
|
val settings = new Settings()
|
|
val settings = new Settings()
|
|
|
settings.usejavacp.value = true
|
|
settings.usejavacp.value = true
|
|
|
|
|
|
|
|
|
|
+ val sparkConf = new SparkConf(true)
|
|
|
|
|
+ .setAppName("Livy Spark shell")
|
|
|
|
|
+
|
|
|
|
|
+ sparkContext = new SparkContext(sparkConf)
|
|
|
|
|
+
|
|
|
sparkIMain = createSparkIMain(classLoader, settings)
|
|
sparkIMain = createSparkIMain(classLoader, settings)
|
|
|
|
|
+ sparkIMain.initializeSynchronous()
|
|
|
|
|
+ sparkIMain.beQuietDuring {
|
|
|
|
|
+ sparkIMain.bind("sc", "org.apache.spark.SparkContext", sparkContext, List("""@transient"""))
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
_state = Interpreter.Idle()
|
|
_state = Interpreter.Idle()
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ private def getMaster(): String = {
|
|
|
|
|
+ sys.props.get("spark.master").getOrElse("local[*]")
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
private def createSparkIMain(classLoader: ClassLoader, settings: Settings) = {
|
|
private def createSparkIMain(classLoader: ClassLoader, settings: Settings) = {
|
|
|
val out = new JPrintWriter(outputStream, true)
|
|
val out = new JPrintWriter(outputStream, true)
|
|
|
val cls = classLoader.loadClass(classOf[SparkIMain].getName)
|
|
val cls = classLoader.loadClass(classOf[SparkIMain].getName)
|