【问题标题】:Best Practice to launch Spark Applications via Web Application?通过 Web 应用程序启动 Spark 应用程序的最佳实践?
【发布时间】:2021-07-01 16:47:19
【问题描述】:

我想通过 Web 应用程序向用户公开我的 Spark 应用程序。

基本上,用户可以决定他想要运行哪个动作并输入一些变量,这些变量需要传递给 spark 应用程序。 例如:用户输入几个字段,然后单击一个按钮,该按钮执行以下“使用参数 min_x、max_x、min_y、max_y 运行 sparkApp1”。

应使用用户提供的参数启动 spark 应用程序。完成后,可能需要 Web 应用程序检索结果(来自 hdfs 或 mongodb)并将其显示给用户。处理时,Web 应用程序应显示 Spark 应用程序的状态。

我的问题:

  • Web 应用程序如何启动 Spark 应用程序?它可能能够从后台的命令行启动它,但可能有更好的方法来做到这一点。
  • Web 应用程序如何访问 Spark 应用程序的当前状态?从 Spark WebUI 的 REST API 获取状态是否可行?

我正在使用 YARN/Mesos(尚不确定)和 MongoDB 运行 Spark 1.6.1 集群。

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    非常基本的答案:

    基本上您可以使用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 ServerSpark 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。

    比较与结论

    我列出了几个选项:

    1. SparkLauncher
    2. Spark REST API
    3. Livy REST 服务器和 Spark 作业服务器
    4. 水圈雾
    5. 阿帕奇托雷

    它们都适用于不同的用例。我可以区分几个类别:

    1. 需要 JAR 文件和作业的工具:Spark Launcher、Spark REST API
    2. 交互式和预打包作业的​​工具:Livy、SJS、Mist
    3. 专注于交互式分析的工具: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 提交用户代码

    【讨论】:

    • 对于 Spark Launcher,我们似乎无法远程启动 spark 作业。我们需要在可以访问集群(边缘节点)的机器上编写启动器作业。如果我错了,请纠正我?
    【解决方案2】:

    这里提到了SparkLauncherT.Gawęda 的例子:

    SparkAppHandle handle = new SparkLauncher()
        .setSparkHome(SPARK_HOME)
        .setJavaHome(JAVA_HOME)
        .setAppResource(SPARK_JOB_JAR_PATH)
        .setMainClass(SPARK_JOB_MAIN_CLASS)
        .addAppArgs("arg1", "arg2")
        .setMaster("yarn-cluster")
        .setConf("spark.dynamicAllocation.enabled", "true")
        .startApplication();
    

    Here 你可以找到一个 Java Web 应用程序示例,其中 Spark 作业捆绑在一个项目中。通过SparkLauncher,您可以获得SparkAppHandle,您可以使用它来获取有关工作状态的信息。如果您需要进度状态,可以使用 Spark rest-api

    http://driverHost:4040/api/v1/applications/[app-id]/jobs
    

    SparkLauncher 需要的唯一依赖项:

    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-launcher_2.10</artifactId>
        <version>2.0.1</version>
    </dependency>
    

    【讨论】:

      【解决方案3】:

      您可以使用 PredictionIO PredictionIO,一个面向开发人员和 ML 工程师的机器学习服务器。 https://github.com/apache/predictionio

      【讨论】:

        【解决方案4】:

        我觉得你可以试试SparkSQL JDBC的方式。

        首先,你应该使用命令./dev/make-distribution.sh --name custom-spark --tgz -PR -Phive -Phive-thriftserver -Pyarn -Dhadoop.version=3.2.0克隆spark源代码和rebuild,不同的jar版本可能会导致兼容的运行时错误。

        其次,安装上面的包并启动thriftserversbin/start-thriftserver.sh --master yarn --deploy-mode client

        那么你可以使用beeline来测试JDBC是否正常工作。

        【讨论】:

          猜你喜欢
          • 2011-02-18
          • 2012-06-06
          • 2015-08-05
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2011-02-16
          • 2011-05-19
          相关资源
          最近更新 更多