001/**
002 * Licensed to the Apache Software Foundation (ASF) under one
003 * or more contributor license agreements.  See the NOTICE file
004 * distributed with this work for additional information
005 * regarding copyright ownership.  The ASF licenses this file
006 * to you under the Apache License, Version 2.0 (the
007 * "License"); you may not use this file except in compliance
008 * with the License.  You may obtain a copy of the License at
009 *
010 *      http://www.apache.org/licenses/LICENSE-2.0
011 *
012 * Unless required by applicable law or agreed to in writing, software
013 * distributed under the License is distributed on an "AS IS" BASIS,
014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
015 * See the License for the specific language governing permissions and
016 * limitations under the License.
017 */
018
019package org.apache.oozie.action.hadoop;
020
021import org.apache.hadoop.conf.Configuration;
022import org.apache.hadoop.fs.FileSystem;
023import org.apache.hadoop.fs.Path;
024import org.apache.hadoop.mapred.JobConf;
025import org.apache.oozie.action.ActionExecutorException;
026import org.apache.oozie.client.WorkflowAction;
027import org.apache.oozie.service.Services;
028import org.apache.oozie.service.SparkConfigurationService;
029import org.jdom.Element;
030import org.jdom.Namespace;
031
032import java.util.ArrayList;
033import java.util.List;
034import java.util.Map;
035
036public class SparkActionExecutor extends JavaActionExecutor {
037    public static final String SPARK_MAIN_CLASS_NAME = "org.apache.oozie.action.hadoop.SparkMain";
038    public static final String TASK_USER_PRECEDENCE = "mapreduce.task.classpath.user.precedence"; // hadoop-2
039    public static final String TASK_USER_CLASSPATH_PRECEDENCE = "mapreduce.user.classpath.first";  // hadoop-1
040    public static final String SPARK_MASTER = "oozie.spark.master";
041    public static final String SPARK_MODE = "oozie.spark.mode";
042    public static final String SPARK_OPTS = "oozie.spark.spark-opts";
043    public static final String SPARK_JOB_NAME = "oozie.spark.name";
044    public static final String SPARK_CLASS = "oozie.spark.class";
045    public static final String SPARK_JAR = "oozie.spark.jar";
046    public static final String MAPRED_CHILD_ENV = "mapred.child.env";
047
048    public SparkActionExecutor() {
049        super("spark");
050    }
051
052    @Override
053    Configuration setupActionConf(Configuration actionConf, Context context, Element actionXml, Path appPath)
054            throws ActionExecutorException {
055        actionConf = super.setupActionConf(actionConf, context, actionXml, appPath);
056        Namespace ns = actionXml.getNamespace();
057
058        String master = actionXml.getChildTextTrim("master", ns);
059        actionConf.set(SPARK_MASTER, master);
060
061        String mode = actionXml.getChildTextTrim("mode", ns);
062        if (mode != null) {
063            actionConf.set(SPARK_MODE, mode);
064        }
065
066        String jobName = actionXml.getChildTextTrim("name", ns);
067        actionConf.set(SPARK_JOB_NAME, jobName);
068
069        String sparkClass = actionXml.getChildTextTrim("class", ns);
070        if (sparkClass != null) {
071            actionConf.set(SPARK_CLASS, sparkClass);
072        }
073
074        String jarLocation = actionXml.getChildTextTrim("jar", ns);
075        actionConf.set(SPARK_JAR, jarLocation);
076
077        StringBuilder sparkOptsSb = new StringBuilder();
078        if (master.startsWith("yarn")) {
079            String resourceManager = actionConf.get(HADOOP_JOB_TRACKER);
080            Map<String, String> sparkConfig = Services.get().get(SparkConfigurationService.class).getSparkConfig(resourceManager);
081            for (Map.Entry<String, String> entry : sparkConfig.entrySet()) {
082                sparkOptsSb.append("--conf ").append(entry.getKey()).append("=").append(entry.getValue()).append(" ");
083            }
084            // CLOUDERA-BUILD: Tell Spark not to localize Hadoop Configs to fix regression on kerberized cluster (CDH-32176)
085            sparkOptsSb.append("--conf spark.yarn.localizeConfig=false ");
086        }
087        String sparkOpts = actionXml.getChildTextTrim("spark-opts", ns);
088        if (sparkOpts != null) {
089            sparkOptsSb.append(sparkOpts);
090        }
091        if (sparkOptsSb.length() > 0) {
092            actionConf.set(SPARK_OPTS, sparkOptsSb.toString().trim());
093        }
094
095        return actionConf;
096    }
097
098    @Override
099    JobConf createLauncherConf(FileSystem actionFs, Context context, WorkflowAction action, Element actionXml,
100                               Configuration actionConf) throws ActionExecutorException {
101
102        JobConf launcherJobConf = super.createLauncherConf(actionFs, context, action, actionXml, actionConf);
103        if (launcherJobConf.get("oozie.launcher." + TASK_USER_PRECEDENCE) == null) {
104            launcherJobConf.set(TASK_USER_PRECEDENCE, "true");
105        }
106        if (launcherJobConf.get("oozie.launcher." + TASK_USER_CLASSPATH_PRECEDENCE) == null) {
107            launcherJobConf.set(TASK_USER_CLASSPATH_PRECEDENCE, "true");
108        }
109        return launcherJobConf;
110    }
111
112    @Override
113    Configuration setupLauncherConf(Configuration conf, Element actionXml, Path appPath, Context context)
114            throws ActionExecutorException {
115        super.setupLauncherConf(conf, actionXml, appPath, context);
116
117        // Set SPARK_HOME environment variable on launcher job
118        // It is needed since pyspark client checks for it.
119        String sparkHome = "SPARK_HOME=.";
120        String mapredChildEnv = conf.get("oozie.launcher." + MAPRED_CHILD_ENV);
121
122        if (mapredChildEnv == null) {
123            conf.set(MAPRED_CHILD_ENV, sparkHome);
124            conf.set("oozie.launcher." + MAPRED_CHILD_ENV, sparkHome);
125        } else if (!mapredChildEnv.contains("SPARK_HOME")) {
126            conf.set(MAPRED_CHILD_ENV, mapredChildEnv + "," + sparkHome);
127            conf.set("oozie.launcher." + MAPRED_CHILD_ENV, mapredChildEnv + "," + sparkHome);
128        }
129        return conf;
130    }
131
132    @Override
133    public List<Class> getLauncherClasses() {
134        List<Class> classes = new ArrayList<Class>();
135        try {
136            classes.add(Class.forName(SPARK_MAIN_CLASS_NAME));
137        } catch (ClassNotFoundException e) {
138            throw new RuntimeException("Class not found", e);
139        }
140        return classes;
141    }
142
143
144    /**
145     * Return the sharelib name for the action.
146     *
147     * @param actionXml
148     * @return returns <code>spark</code>.
149     */
150    @Override
151    protected String getDefaultShareLibName(Element actionXml) {
152        return "spark";
153    }
154
155    @Override
156    protected String getLauncherMain(Configuration launcherConf, Element actionXml) {
157        return launcherConf.get(LauncherMapper.CONF_OOZIE_ACTION_MAIN_CLASS, SPARK_MAIN_CLASS_NAME);
158    }
159}