|
@@ -5,7 +5,13 @@ import com.cloudera.hue.livy.Logging
|
|
|
import scala.collection.JavaConversions._
|
|
import scala.collection.JavaConversions._
|
|
|
import scala.collection.mutable.ArrayBuffer
|
|
import scala.collection.mutable.ArrayBuffer
|
|
|
|
|
|
|
|
-class SparkProcessBuilder extends Logging {
|
|
|
|
|
|
|
+object SparkSubmitProcessBuilder {
|
|
|
|
|
+ def apply(): SparkSubmitProcessBuilder = {
|
|
|
|
|
+ new SparkSubmitProcessBuilder()
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+class SparkSubmitProcessBuilder extends Logging {
|
|
|
|
|
|
|
|
private[this] var _executable = "spark-submit"
|
|
private[this] var _executable = "spark-submit"
|
|
|
private[this] var _master: Option[String] = None
|
|
private[this] var _master: Option[String] = None
|
|
@@ -33,142 +39,150 @@ class SparkProcessBuilder extends Logging {
|
|
|
private[this] var _redirectError: Option[ProcessBuilder.Redirect] = None
|
|
private[this] var _redirectError: Option[ProcessBuilder.Redirect] = None
|
|
|
private[this] var _redirectErrorStream: Option[Boolean] = None
|
|
private[this] var _redirectErrorStream: Option[Boolean] = None
|
|
|
|
|
|
|
|
- def executable(executable: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def executable(executable: String): SparkSubmitProcessBuilder = {
|
|
|
_executable = executable
|
|
_executable = executable
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def master(masterUrl: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def master(masterUrl: String): SparkSubmitProcessBuilder = {
|
|
|
_master = Some(masterUrl)
|
|
_master = Some(masterUrl)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def deployMode(deployMode: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def deployMode(deployMode: String): SparkSubmitProcessBuilder = {
|
|
|
_deployMode = Some(deployMode)
|
|
_deployMode = Some(deployMode)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def className(className: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def className(className: String): SparkSubmitProcessBuilder = {
|
|
|
_className = Some(className)
|
|
_className = Some(className)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def name(name: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def name(name: String): SparkSubmitProcessBuilder = {
|
|
|
_name = Some(name)
|
|
_name = Some(name)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def jar(jar: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def jar(jar: String): SparkSubmitProcessBuilder = {
|
|
|
this._jars += jar
|
|
this._jars += jar
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def jars(jars: Traversable[String]): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def jars(jars: Traversable[String]): SparkSubmitProcessBuilder = {
|
|
|
this._jars ++= jars
|
|
this._jars ++= jars
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def pyFile(pyFile: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def pyFile(pyFile: String): SparkSubmitProcessBuilder = {
|
|
|
this._pyFiles += pyFile
|
|
this._pyFiles += pyFile
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def pyFiles(pyFiles: Traversable[String]): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def pyFiles(pyFiles: Traversable[String]): SparkSubmitProcessBuilder = {
|
|
|
this._pyFiles ++= pyFiles
|
|
this._pyFiles ++= pyFiles
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def file(file: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def file(file: String): SparkSubmitProcessBuilder = {
|
|
|
this._files += file
|
|
this._files += file
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def files(files: Traversable[String]): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def files(files: Traversable[String]): SparkSubmitProcessBuilder = {
|
|
|
this._files ++= files
|
|
this._files ++= files
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def conf(key: String, value: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def conf(key: String, value: String): SparkSubmitProcessBuilder = {
|
|
|
this._conf += ((key, value))
|
|
this._conf += ((key, value))
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def conf(conf: Traversable[(String, String)]): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def conf(conf: Traversable[(String, String)]): SparkSubmitProcessBuilder = {
|
|
|
this._conf ++= conf
|
|
this._conf ++= conf
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def driverMemory(driverMemory: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def driverMemory(driverMemory: String): SparkSubmitProcessBuilder = {
|
|
|
_driverMemory = Some(driverMemory)
|
|
_driverMemory = Some(driverMemory)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def driverJavaOptions(driverJavaOptions: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def driverJavaOptions(driverJavaOptions: String): SparkSubmitProcessBuilder = {
|
|
|
_driverJavaOptions = Some(driverJavaOptions)
|
|
_driverJavaOptions = Some(driverJavaOptions)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def driverClassPath(classPath: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def driverClassPath(classPath: String): SparkSubmitProcessBuilder = {
|
|
|
_driverClassPath += classPath
|
|
_driverClassPath += classPath
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def driverClassPaths(classPaths: Traversable[String]): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def driverClassPaths(classPaths: Traversable[String]): SparkSubmitProcessBuilder = {
|
|
|
_driverClassPath ++= classPaths
|
|
_driverClassPath ++= classPaths
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def executorMemory(executorMemory: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def executorMemory(executorMemory: String): SparkSubmitProcessBuilder = {
|
|
|
_executorMemory = Some(executorMemory)
|
|
_executorMemory = Some(executorMemory)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def proxyUser(proxyUser: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def proxyUser(proxyUser: String): SparkSubmitProcessBuilder = {
|
|
|
_proxyUser = Some(proxyUser)
|
|
_proxyUser = Some(proxyUser)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def driverCores(driverCores: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def driverCores(driverCores: Int): SparkSubmitProcessBuilder = {
|
|
|
|
|
+ this.driverCores(driverCores.toString)
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ def driverCores(driverCores: String): SparkSubmitProcessBuilder = {
|
|
|
_driverCores = Some(driverCores)
|
|
_driverCores = Some(driverCores)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def executorCores(executorCores: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def executorCores(executorCores: Int): SparkSubmitProcessBuilder = {
|
|
|
|
|
+ this.executorCores(executorCores.toString)
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ def executorCores(executorCores: String): SparkSubmitProcessBuilder = {
|
|
|
_executorCores = Some(executorCores)
|
|
_executorCores = Some(executorCores)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def queue(queue: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def queue(queue: String): SparkSubmitProcessBuilder = {
|
|
|
_queue = Some(queue)
|
|
_queue = Some(queue)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def archive(archive: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def archive(archive: String): SparkSubmitProcessBuilder = {
|
|
|
_archives += archive
|
|
_archives += archive
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def archives(archives: Traversable[String]): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def archives(archives: Traversable[String]): SparkSubmitProcessBuilder = {
|
|
|
_archives ++= archives
|
|
_archives ++= archives
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def env(key: String, value: String): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def env(key: String, value: String): SparkSubmitProcessBuilder = {
|
|
|
_env += ((key, value))
|
|
_env += ((key, value))
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def redirectOutput(redirect: ProcessBuilder.Redirect): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def redirectOutput(redirect: ProcessBuilder.Redirect): SparkSubmitProcessBuilder = {
|
|
|
_redirectOutput = Some(redirect)
|
|
_redirectOutput = Some(redirect)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def redirectError(redirect: ProcessBuilder.Redirect): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def redirectError(redirect: ProcessBuilder.Redirect): SparkSubmitProcessBuilder = {
|
|
|
_redirectError = Some(redirect)
|
|
_redirectError = Some(redirect)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def redirectErrorStream(redirect: Boolean): SparkProcessBuilder = {
|
|
|
|
|
|
|
+ def redirectErrorStream(redirect: Boolean): SparkSubmitProcessBuilder = {
|
|
|
_redirectErrorStream = Some(redirect)
|
|
_redirectErrorStream = Some(redirect)
|
|
|
this
|
|
this
|
|
|
}
|
|
}
|