非常基本的答案:
基本上您可以使用SparkLauncher 类来启动 Spark 应用程序并添加一些侦听器来查看进度。
不过,您可能对 Livy 服务器感兴趣,它是用于 Spark 作业的 RESTful 服务器。据我所知,Zeppelin 正在使用 Livy 提交作业和检索状态。
您也可以使用 Spark REST 接口来检查状态,这样信息会更准确。 Here 有一个如何通过 REST API 提交作业的示例
您有 3 个选项,答案是 - 自己检查 ;) 这很大程度上取决于您的项目和要求。两个主要选项:
- SparkLauncher + Spark REST 接口
- Livy 服务器
应该对你有好处,你必须检查一下在你的项目中使用什么更容易和更好
扩展答案
您可以通过不同的方式从您的应用程序中使用 Spark,具体取决于您的需求和偏好。
SparkLauncher
SparkLauncher 是来自spark-launcher 工件的类。它用于启动已经准备好的 Spark 作业,就像从 Spark Submit 一样。
典型用法是:
1) 使用您的 Spark 作业构建项目并将 JAR 文件复制到所有节点
2) 从您的客户端应用程序,即 Web 应用程序,创建指向准备好的 JAR 文件的 SparkLauncher
SparkAppHandle handle = new SparkLauncher()
.setSparkHome(SPARK_HOME)
.setJavaHome(JAVA_HOME)
.setAppResource(pathToJARFile)
.setMainClass(MainClassFromJarWithJob)
.setMaster("MasterAddress
.startApplication();
// or: .launch().waitFor()
startApplication 创建 SparkAppHandle,它允许您添加侦听器并停止应用程序。它还提供了getAppId的可能性。
SparkLauncher 应与 Spark REST API 一起使用。您可以查询http://driverNode:4040/api/v1/applications/*ResultFromGetAppId*/jobs,您将获得有关应用程序当前状态的信息。
Spark REST API
也可以通过 RESTful API 直接提交 Spark 作业。用法与SparkLauncher 非常相似,但它是以纯RESTful 方式完成的。
示例请求 - 这篇文章的学分:
curl -X POST http://spark-master-host:6066/v1/submissions/create --header "Content-Type:application/json;charset=UTF-8" --data '{
"action" : "CreateSubmissionRequest",
"appArgs" : [ "myAppArgument1" ],
"appResource" : "hdfs:///filepath/spark-job-1.0.jar",
"clientSparkVersion" : "1.5.0",
"environmentVariables" : {
"SPARK_ENV_LOADED" : "1"
},
"mainClass" : "spark.ExampleJobInPreparedJar",
"sparkProperties" : {
"spark.jars" : "hdfs:///filepath/spark-job-1.0.jar",
"spark.driver.supervise" : "false",
"spark.app.name" : "ExampleJobInPreparedJar",
"spark.eventLog.enabled": "true",
"spark.submit.deployMode" : "cluster",
"spark.master" : "spark://spark-cluster-ip:6066"
}
}'
此命令会将ExampleJobInPreparedJar 类中的作业提交到具有给定 Spark Master 的集群。在响应中,您将有submissionId 字段,这将有助于检查应用程序的状态 - 只需调用另一个服务:curl http://spark-cluster-ip:6066/v1/submissions/status/submissionIdFromResponse。就是这样,没有更多的代码
Livy REST 服务器和 Spark 作业服务器
Livy REST Server 和 Spark Job Server 是 RESTful 应用程序,允许您通过 RESTful Web Service 提交作业。这两者和 Spark 的 REST 接口之间的一个主要区别是 Livy 和 SJS 不需要提前准备作业并将其打包到 JAR 文件中。您只是提交将在 Spark 中执行的代码。
用法很简单。代码取自 Livy 存储库,但为了提高可读性做了一些删减
1) 案例1:提交作业,放在本地机器上
// creating client
LivyClient client = new LivyClientBuilder()
.setURI(new URI(livyUrl))
.build();
try {
// sending and submitting JAR file
client.uploadJar(new File(piJar)).get();
// PiJob is a class that implements Livy's Job
double pi = client.submit(new PiJob(samples)).get();
} finally {
client.stop(true);
}
2) 案例 2:动态作业创建和执行
// example in Python. Data contains code in Scala, that will be executed in Spark
data = {
'code': textwrap.dedent("""\
val NUM_SAMPLES = 100000;
val count = sc.parallelize(1 to NUM_SAMPLES).map { i =>
val x = Math.random();
val y = Math.random();
if (x*x + y*y < 1) 1 else 0
}.reduce(_ + _);
println(\"Pi is roughly \" + 4.0 * count / NUM_SAMPLES)
""")
}
r = requests.post(statements_url, data=json.dumps(data), headers=headers)
pprint.pprint(r.json())
如您所见,预编译作业和对 Spark 的即席查询都是可能的。
水圈雾
另一个 Spark 即服务应用程序。 Mist 非常简单,类似于 Livy 和 Spark Job Server。
用法非常相似
1) 创建作业文件:
import io.hydrosphere.mist.MistJob
object MyCoolMistJob extends MistJob {
def doStuff(parameters: Map[String, Any]): Map[String, Any] = {
val rdd = context.parallelize()
...
return result.asInstance[Map[String, Any]]
}
}
2) 将作业文件打包成 JAR
3) 向 Mist 发送请求:
curl --header "Content-Type: application/json" -X POST http://mist_http_host:mist_http_port/jobs --data '{"path": "/path_to_jar/mist_examples.jar", "className": "SimpleContext$", "parameters": {"digits": [1, 2, 3, 4, 5, 6, 7, 8, 9, 0]}, "namespace": "foo"}'
我可以在 Mist 中看到一个强大的功能,那就是它通过 MQTT 对流式作业提供开箱即用的支持。
阿帕奇托雷
Apache Toree 旨在为 Spark 启用简单的交互式分析。它不需要构建任何 JAR。它通过 IPython 协议工作,但不仅支持 Python。
当前文档侧重于 Jupyter notebook 支持,但也有 REST 样式的 API。
比较与结论
我列出了几个选项:
- SparkLauncher
- Spark REST API
- Livy REST 服务器和 Spark 作业服务器
- 水圈雾
- 阿帕奇托雷
它们都适用于不同的用例。我可以区分几个类别:
- 需要 JAR 文件和作业的工具:Spark Launcher、Spark REST API
- 交互式和预打包作业的工具:Livy、SJS、Mist
- 专注于交互式分析的工具:Toree(但可能对预先打包的作业提供一些支持;目前未发布任何文档)
SparkLauncher 非常简单,是 Spark 项目的一部分。您正在用纯代码编写作业配置,因此它比 JSON 对象更容易构建。
对于完全 RESTful 风格的提交,请考虑 Spark REST API、Livy、SJS 和 Mist。其中三个是稳定的项目,有一些生产用例。 REST API 还需要预先打包作业,而 Livy 和 SJS 不需要。但是请记住,Spark REST API 在每个 Spark 发行版中都是默认的,而 Livy/SJS 不是。我对 Mist 了解不多,但是 - 一段时间后 - 它应该是集成所有类型的 Spark 作业的非常好的工具。
Toree 专注于交互式工作。它仍处于孵化阶段,但即使现在你也可以检查它的可能性。
既然有内置的 REST API,为什么还要使用自定义的附加 REST 服务?像 Livy 这样的 SaaS 是 Spark 的一个入口点。它管理 Spark 上下文,并且只在一个节点上,而不是在集群之外的其他地方。它们还支持交互式分析。 Apache Zeppelin 使用 Livy 向 Spark 提交用户代码