当前位置: 代码迷 >> java >> 如果在SparkAction中使用PySpark,则Oozie作业将无法运行
  详细解决方案

如果在SparkAction中使用PySpark,则Oozie作业将无法运行

热度:110   发布时间:2023-07-26 14:31:22.0

我在Oozie中遇到了几个SparkAction作业示例,其中大多数都是Java语言。 我做了一些编辑,然后在Cloudera CDH Quickstart 5.4.0(Spark版本1.4.0)中运行示例。

工作流程

<workflow-app xmlns='uri:oozie:workflow:0.5' name='SparkFileCopy'>
    <start to='spark-node' />

    <action name='spark-node'>
        <spark xmlns="uri:oozie:spark-action:0.1">
            <job-tracker>${jobTracker}</job-tracker>
            <name-node>${nameNode}</name-node>
            <prepare>
                <delete path="${nameNode}/user/${wf:user()}/${examplesRoot}/output-data/spark"/>
            </prepare>
            <master>${master}</master>
        <mode>${mode}</mode>    
            <name>Spark-FileCopy</name>
            <class>org.apache.oozie.example.SparkFileCopy</class>
            <jar>${nameNode}/user/${wf:user()}/${examplesRoot}/apps/spark/lib/oozie-examples.jar</jar>
            <arg>${nameNode}/user/${wf:user()}/${examplesRoot}/input-data/text/data.txt</arg>
            <arg>${nameNode}/user/${wf:user()}/${examplesRoot}/output-data/spark</arg>
        </spark>
        <ok to="end" />
        <error to="fail" />
    </action>

    <kill name="fail">
        <message>Workflow failed, error
            message[${wf:errorMessage(wf:lastErrorNode())}]
        </message>
    </kill>
    <end name='end' />
</workflow-app>

job.properties

nameNode=hdfs://quickstart.cloudera:8020
jobTracker=quickstart.cloudera:8032
master=local[2]
mode=client
examplesRoot=examples
oozie.use.system.libpath=true
oozie.wf.application.path=${nameNode}/user/${user.name}/${examplesRoot}/apps/spark

Oozie工作流示例(使用Java)能够完成并完成其任务。

但是,我已经使用Python / PySpark编写了spark-submit作业。 我尝试删除<class>和罐子

<jar>my_pyspark_job.py</jar>

但是当我尝试运行Oozie-Spark作业时,日志中出现错误:

Launcher ERROR, reason: Main class [org.apache.oozie.action.hadoop.SparkMain], exit code [2]

我想知道如果我使用的是Python / PySpark,应该在<class><jar>标记中放置什么?

我也为oozie中的火花动作而苦苦挣扎。 我正确设置了sharelib,并尝试使用<spark-opts> </spark-opts>标记内的--jars选项传递适当的jar,但无济于事。

我总是总是遇到一些错误或其他错误。 我最能做的就是通过spark-action在本地模式下运行所有??java / python spark作业。

但是,我使用shell动作以所有执行模式在oozie中运行了所有spark作业。 shell动作的主要问题是将shell作业部署为“ yarn”用户。 如果您碰巧从不是纱线的用户帐户部署oozie spark作业,最终将出现“权限被拒绝”错误(因为用户将无法访问复制到/user/yarn/.SparkStaging中的spark程序罐。目录)。 解决此问题的方法是将HADOOP_USER_NAME环境变量设置为用于部署oozie工作流的用户帐户名。

下面是说明此配置的工作流程。 我从ambari-qa用户部署oozie工作流。

<workflow-app xmlns="uri:oozie:workflow:0.4" name="sparkjob">
    <start to="spark-shell-node"/>
    <action name="spark-shell-node">
        <shell xmlns="uri:oozie:shell-action:0.2">
            <job-tracker>${jobTracker}</job-tracker>
            <name-node>${nameNode}</name-node>
            <configuration>
                <property>
                    <name>oozie.launcher.mapred.job.queue.name</name>
                    <value>launcher2</value>
                </property>
                <property>
                    <name>mapred.job.queue.name</name>
                    <value>default</value>
                </property>
                <property>
                    <name>oozie.hive.defaults</name>
                    <value>/user/ambari-qa/sparkActionPython/hive-site.xml</value>
                </property>
            </configuration>
            <exec>/usr/hdp/current/spark-client/bin/spark-submit</exec>
            <argument>--master</argument>
            <argument>yarn-cluster</argument>
            <argument>wordcount.py</argument>
            <env-var>HADOOP_USER_NAME=ambari-qa</env-var>
            <file>/user/ambari-qa/sparkActionPython/wordcount.py#wordcount.py</file>
            <capture-output/>
        </shell>
        <ok to="end"/>
        <error to="spark-fail"/>
    </action>
    <kill name="spark-fail">
        <message>Shell action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message>
    </kill>
    <end name="end"/>
</workflow-app>

希望这可以帮助!

您应该尝试配置Oozie Spark操作以将所需的文件带到本地。 您可以使用文件标记来使其:

<spark xmlns="uri:oozie:spark-action:0.1">
        <job-tracker>${resourceManager}</job-tracker>
        <name-node>${nameNode}</name-node>
        <master>local[2]</master>
        <mode>client</mode>
        <name>${name}</name>
        <jar>my_pyspark_job.py</jar>
        <file>{path to your file on hdfs}/my_pyspark_job.py#my_pyspark_job.py</file>
    </spark>

说明:在YARN容器中运行的Oozie操作,由YARN在具有可用资源的节点上分配。 在运行该动作(实际上是一个“驱动程序”代码)之前,它将所有需要的文件(例如jar)本地复制到该节点到分配给YARN容器的文件夹中,以放置其资源。 因此,通过向oozie动作添加标签,您可以“告诉”您的oozie动作,以将my_pyspark_job.py本地带到执行节点。

就我而言,我想运行一个bash脚本(run-hive-partitioner.bash),该脚本将运行python代码(hive-generic-partitioner.py),因此我需要在该节点上本地可访问的所有文件:

<action name="repair_hive_partitions">
  <shell xmlns="uri:oozie:shell-action:0.1">
    <job-tracker>${jobTracker}</job-tracker>
    <name-node>${nameNode}</name-node>
    <exec>${appPath}/run-hive-partitioner.bash</exec>
         <argument>${db}</argument>
         <argument>${tables}</argument>
         <argument>${base_working_dir}</argument>
    <file>${appPath}/run-hive-partitioner.bash#run-hive-partitioner.bash</file>
    <file>${appPath}/hive-generic-partitioner.py#hive-generic-partitioner.py</file>
     <file>${appPath}/util.py#util.py</file>     
  </shell>
  <ok to="end"/>
  <error to="kill"/>
</action>

其中$ {appPath}是hdfs://ci-base.com:8020 / app / oozie / util / wf-repair_hive_partitions

所以这就是我的工作:

Files in current dir:/hadoop/yarn/local/usercache/hdfs/appcache/application_1440506439954_3906/container_1440506439954_3906_01_000002/

======================
File: hive-generic-partitioner.py
File: util.py
File: run-hive-partitioner.bash
...
File: job.xml
File: json-simple-1.1.jar
File: oozie-sharelib-oozie-4.1.0.2.2.4.2-2.jar
File: launch_container.sh
File: oozie-hadoop-utils-2.6.0.2.2.4.2-2.oozie-4.1.0.2.2.4.2-2.jar

如您所见,oozie(或我认为实际上是毛线)将所有需要的文件本地发送到temp文件夹,现在它可以运行它了。

我能够“解决”此问题,尽管它导致了另一个问题。 尽管如此,我仍然会发布它。

在Oozie容器日志的stderr中,它显示:

Error: Only local python files are supported

我在找到了解决方案

这是我的初始工作流程.xml:

    <spark xmlns="uri:oozie:spark-action:0.1">
        <job-tracker>${resourceManager}</job-tracker>
        <name-node>${nameNode}</name-node>
        <master>local[2]</master>
        <mode>client</mode>
        <name>${name}</name>
        <jar>my_pyspark_job.py</jar>
    </spark>

我最初所做的是将希望作为火花提交作业运行的Python脚本复制到HDFS。 事实证明,它期望本地文件系统中有.py脚本,因此我要做的是引用脚本的绝对本地文件系统。

<jar>/<absolute-local-path>/my_pyspark_job.py</jar>

我们遇到了同样的错误。 如果您尝试将oozie.wf.application.path/lib -assembly jar从'/path/to/spark-install/lib/spark-assembly*.jar'(取决于分发) oozie.wf.application.path/lib到应用程序旁边的oozie.wf.application.path/lib目录中罐子应该工作。