|
|
@@ -18,15 +18,15 @@
|
|
|
|
|
|
package com.cloudera.hue.livy.spark
|
|
|
|
|
|
-import com.cloudera.hue.livy.spark.SparkSubmitProcessBuilder.{RelativePath, AbsolutePath, Path}
|
|
|
+import com.cloudera.hue.livy.spark.SparkProcessBuilder.{RelativePath, AbsolutePath, Path}
|
|
|
import com.cloudera.hue.livy.{LivyConf, Logging}
|
|
|
|
|
|
import scala.collection.JavaConversions._
|
|
|
import scala.collection.mutable.ArrayBuffer
|
|
|
|
|
|
-object SparkSubmitProcessBuilder {
|
|
|
- def apply(livyConf: LivyConf): SparkSubmitProcessBuilder = {
|
|
|
- new SparkSubmitProcessBuilder(livyConf)
|
|
|
+object SparkProcessBuilder {
|
|
|
+ def apply(livyConf: LivyConf): SparkProcessBuilder = {
|
|
|
+ new SparkProcessBuilder(livyConf)
|
|
|
}
|
|
|
|
|
|
/**
|
|
|
@@ -38,7 +38,7 @@ object SparkSubmitProcessBuilder {
|
|
|
case class RelativePath(path: String) extends Path
|
|
|
}
|
|
|
|
|
|
-class SparkSubmitProcessBuilder(livyConf: LivyConf) extends Logging {
|
|
|
+class SparkProcessBuilder(livyConf: LivyConf) extends Logging {
|
|
|
private[this] val fsRoot = livyConf.filesystemRoot()
|
|
|
|
|
|
private[this] var _executable: Path = AbsolutePath(livyConf.sparkSubmit())
|
|
|
@@ -67,160 +67,160 @@ class SparkSubmitProcessBuilder(livyConf: LivyConf) extends Logging {
|
|
|
private[this] var _redirectError: Option[ProcessBuilder.Redirect] = None
|
|
|
private[this] var _redirectErrorStream: Option[Boolean] = None
|
|
|
|
|
|
- def executable(executable: Path): SparkSubmitProcessBuilder = {
|
|
|
+ def executable(executable: Path): SparkProcessBuilder = {
|
|
|
_executable = executable
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def master(masterUrl: String): SparkSubmitProcessBuilder = {
|
|
|
+ def master(masterUrl: String): SparkProcessBuilder = {
|
|
|
_master = Some(masterUrl)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def deployMode(deployMode: String): SparkSubmitProcessBuilder = {
|
|
|
+ def deployMode(deployMode: String): SparkProcessBuilder = {
|
|
|
_deployMode = Some(deployMode)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def className(className: String): SparkSubmitProcessBuilder = {
|
|
|
+ def className(className: String): SparkProcessBuilder = {
|
|
|
_className = Some(className)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def name(name: String): SparkSubmitProcessBuilder = {
|
|
|
+ def name(name: String): SparkProcessBuilder = {
|
|
|
_name = Some(name)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def jar(jar: Path): SparkSubmitProcessBuilder = {
|
|
|
+ def jar(jar: Path): SparkProcessBuilder = {
|
|
|
this._jars += jar
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def jars(jars: Traversable[Path]): SparkSubmitProcessBuilder = {
|
|
|
+ def jars(jars: Traversable[Path]): SparkProcessBuilder = {
|
|
|
this._jars ++= jars
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def pyFile(pyFile: Path): SparkSubmitProcessBuilder = {
|
|
|
+ def pyFile(pyFile: Path): SparkProcessBuilder = {
|
|
|
this._pyFiles += pyFile
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def pyFiles(pyFiles: Traversable[Path]): SparkSubmitProcessBuilder = {
|
|
|
+ def pyFiles(pyFiles: Traversable[Path]): SparkProcessBuilder = {
|
|
|
this._pyFiles ++= pyFiles
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def file(file: Path): SparkSubmitProcessBuilder = {
|
|
|
+ def file(file: Path): SparkProcessBuilder = {
|
|
|
this._files += file
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def files(files: Traversable[Path]): SparkSubmitProcessBuilder = {
|
|
|
+ def files(files: Traversable[Path]): SparkProcessBuilder = {
|
|
|
this._files ++= files
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def conf(key: String, value: String): SparkSubmitProcessBuilder = {
|
|
|
+ def conf(key: String, value: String): SparkProcessBuilder = {
|
|
|
this._conf += ((key, value))
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def conf(conf: Traversable[(String, String)]): SparkSubmitProcessBuilder = {
|
|
|
+ def conf(conf: Traversable[(String, String)]): SparkProcessBuilder = {
|
|
|
this._conf ++= conf
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def driverMemory(driverMemory: String): SparkSubmitProcessBuilder = {
|
|
|
+ def driverMemory(driverMemory: String): SparkProcessBuilder = {
|
|
|
_driverMemory = Some(driverMemory)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def driverJavaOptions(driverJavaOptions: String): SparkSubmitProcessBuilder = {
|
|
|
+ def driverJavaOptions(driverJavaOptions: String): SparkProcessBuilder = {
|
|
|
_driverJavaOptions = Some(driverJavaOptions)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def driverClassPath(classPath: String): SparkSubmitProcessBuilder = {
|
|
|
+ def driverClassPath(classPath: String): SparkProcessBuilder = {
|
|
|
_driverClassPath += classPath
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def driverClassPaths(classPaths: Traversable[String]): SparkSubmitProcessBuilder = {
|
|
|
+ def driverClassPaths(classPaths: Traversable[String]): SparkProcessBuilder = {
|
|
|
_driverClassPath ++= classPaths
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def executorMemory(executorMemory: String): SparkSubmitProcessBuilder = {
|
|
|
+ def executorMemory(executorMemory: String): SparkProcessBuilder = {
|
|
|
_executorMemory = Some(executorMemory)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def proxyUser(proxyUser: String): SparkSubmitProcessBuilder = {
|
|
|
+ def proxyUser(proxyUser: String): SparkProcessBuilder = {
|
|
|
_proxyUser = Some(proxyUser)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def driverCores(driverCores: Int): SparkSubmitProcessBuilder = {
|
|
|
+ def driverCores(driverCores: Int): SparkProcessBuilder = {
|
|
|
this.driverCores(driverCores.toString)
|
|
|
}
|
|
|
|
|
|
- def driverCores(driverCores: String): SparkSubmitProcessBuilder = {
|
|
|
+ def driverCores(driverCores: String): SparkProcessBuilder = {
|
|
|
_driverCores = Some(driverCores)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def executorCores(executorCores: Int): SparkSubmitProcessBuilder = {
|
|
|
+ def executorCores(executorCores: Int): SparkProcessBuilder = {
|
|
|
this.executorCores(executorCores.toString)
|
|
|
}
|
|
|
|
|
|
- def executorCores(executorCores: String): SparkSubmitProcessBuilder = {
|
|
|
+ def executorCores(executorCores: String): SparkProcessBuilder = {
|
|
|
_executorCores = Some(executorCores)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
|
|
|
- def numExecutors(numExecutors: Int): SparkSubmitProcessBuilder = {
|
|
|
+ def numExecutors(numExecutors: Int): SparkProcessBuilder = {
|
|
|
this.numExecutors(numExecutors.toString)
|
|
|
}
|
|
|
|
|
|
- def numExecutors(numExecutors: String): SparkSubmitProcessBuilder = {
|
|
|
+ def numExecutors(numExecutors: String): SparkProcessBuilder = {
|
|
|
_numExecutors = Some(numExecutors)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def queue(queue: String): SparkSubmitProcessBuilder = {
|
|
|
+ def queue(queue: String): SparkProcessBuilder = {
|
|
|
_queue = Some(queue)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def archive(archive: Path): SparkSubmitProcessBuilder = {
|
|
|
+ def archive(archive: Path): SparkProcessBuilder = {
|
|
|
_archives += archive
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def archives(archives: Traversable[Path]): SparkSubmitProcessBuilder = {
|
|
|
+ def archives(archives: Traversable[Path]): SparkProcessBuilder = {
|
|
|
archives.foreach(archive)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def env(key: String, value: String): SparkSubmitProcessBuilder = {
|
|
|
+ def env(key: String, value: String): SparkProcessBuilder = {
|
|
|
_env += ((key, value))
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def redirectOutput(redirect: ProcessBuilder.Redirect): SparkSubmitProcessBuilder = {
|
|
|
+ def redirectOutput(redirect: ProcessBuilder.Redirect): SparkProcessBuilder = {
|
|
|
_redirectOutput = Some(redirect)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def redirectError(redirect: ProcessBuilder.Redirect): SparkSubmitProcessBuilder = {
|
|
|
+ def redirectError(redirect: ProcessBuilder.Redirect): SparkProcessBuilder = {
|
|
|
_redirectError = Some(redirect)
|
|
|
this
|
|
|
}
|
|
|
|
|
|
- def redirectErrorStream(redirect: Boolean): SparkSubmitProcessBuilder = {
|
|
|
+ def redirectErrorStream(redirect: Boolean): SparkProcessBuilder = {
|
|
|
_redirectErrorStream = Some(redirect)
|
|
|
this
|
|
|
}
|