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 com.google.common.base.Strings;
022import org.apache.hadoop.conf.Configuration;
023import org.apache.hadoop.fs.FileSystem;
024import org.apache.hadoop.fs.Path;
025import org.apache.hadoop.mapred.JobConf;
026import org.apache.oozie.action.ActionExecutorException;
027import org.apache.oozie.client.WorkflowAction;
028import org.apache.oozie.service.ConfigurationService;
029import org.apache.oozie.service.Services;
030import org.apache.oozie.service.SparkConfigurationService;
031import org.jdom.Element;
032import org.jdom.Namespace;
033
034import java.io.IOException;
035import java.io.StringWriter;
036import java.util.ArrayList;
037import java.util.List;
038import java.util.Properties;
039
040public class SparkActionExecutor extends JavaActionExecutor {
041    public static final String SPARK_MAIN_CLASS_NAME = "org.apache.oozie.action.hadoop.SparkMain";
042    public static final String TASK_USER_PRECEDENCE = "mapreduce.task.classpath.user.precedence"; // hadoop-2
043    public static final String TASK_USER_CLASSPATH_PRECEDENCE = "mapreduce.user.classpath.first";  // hadoop-1
044    public static final String SPARK_MASTER = "oozie.spark.master";
045    public static final String SPARK_MODE = "oozie.spark.mode";
046    public static final String SPARK_OPTS = "oozie.spark.spark-opts";
047    public static final String SPARK_DEFAULT_OPTS = "oozie.spark.spark-default-opts";
048    public static final String SPARK_JOB_NAME = "oozie.spark.name";
049    public static final String SPARK_CLASS = "oozie.spark.class";
050    public static final String SPARK_JAR = "oozie.spark.jar";
051    public static final String MAPRED_CHILD_ENV = "mapred.child.env";
052    private static final String CONF_OOZIE_SPARK_SETUP_HADOOP_CONF_DIR = "oozie.action.spark.setup.hadoop.conf.dir";
053
054    public SparkActionExecutor() {
055        super("spark");
056    }
057
058    @Override
059    Configuration setupActionConf(Configuration actionConf, Context context, Element actionXml, Path appPath)
060            throws ActionExecutorException {
061        actionConf = super.setupActionConf(actionConf, context, actionXml, appPath);
062        Namespace ns = actionXml.getNamespace();
063
064        String master = actionXml.getChildTextTrim("master", ns);
065        actionConf.set(SPARK_MASTER, master);
066
067        String mode = actionXml.getChildTextTrim("mode", ns);
068        if (mode != null) {
069            actionConf.set(SPARK_MODE, mode);
070        }
071
072        String jobName = actionXml.getChildTextTrim("name", ns);
073        actionConf.set(SPARK_JOB_NAME, jobName);
074
075        String sparkClass = actionXml.getChildTextTrim("class", ns);
076        if (sparkClass != null) {
077            actionConf.set(SPARK_CLASS, sparkClass);
078        }
079
080        String jarLocation = actionXml.getChildTextTrim("jar", ns);
081        actionConf.set(SPARK_JAR, jarLocation);
082
083        if (master.startsWith("yarn")) {
084            String resourceManager = actionConf.get(HADOOP_JOB_TRACKER);
085            Properties sparkConfig =
086                    Services.get().get(SparkConfigurationService.class).getSparkConfig(resourceManager);
087            if (!sparkConfig.isEmpty()) {
088                try (final StringWriter sw = new StringWriter()) {
089                    sparkConfig.store(sw, "Generated by Oozie server SparkActionExecutor");
090                    actionConf.set(SPARK_DEFAULT_OPTS, sw.toString());
091                } catch (IOException e) {
092                    LOG.warn("Could not propagate Spark default configuration!", e);
093                }
094            }
095        }
096        String sparkOpts = actionXml.getChildTextTrim("spark-opts", ns);
097        if (!Strings.isNullOrEmpty(sparkOpts)) {
098            // CLOUDERA-BUILD: Tell Spark not to localize Hadoop Configs to fix regression on kerberized cluster (CDH-32176)
099            sparkOpts +=  master.startsWith("yarn") ? " --conf spark.yarn.localizeConfig=false" : "";
100            actionConf.set(SPARK_OPTS, sparkOpts.toString().trim());
101        } else if (master.startsWith("yarn")){
102            // CLOUDERA-BUILD: Tell Spark not to localize Hadoop Configs to fix regression on kerberized cluster (CDH-32176)
103            actionConf.set(SPARK_OPTS, "--conf spark.yarn.localizeConfig=false");
104        }
105
106        // Setting if SparkMain should setup hadoop config *-site.xml
107        boolean setupHadoopConf = actionConf.getBoolean(CONF_OOZIE_SPARK_SETUP_HADOOP_CONF_DIR,
108                ConfigurationService.getBoolean(CONF_OOZIE_SPARK_SETUP_HADOOP_CONF_DIR));
109        actionConf.setBoolean(CONF_OOZIE_SPARK_SETUP_HADOOP_CONF_DIR, setupHadoopConf);
110        return actionConf;
111    }
112
113    @Override
114    JobConf createLauncherConf(FileSystem actionFs, Context context, WorkflowAction action, Element actionXml,
115                               Configuration actionConf) throws ActionExecutorException {
116
117        JobConf launcherJobConf = super.createLauncherConf(actionFs, context, action, actionXml, actionConf);
118        if (launcherJobConf.get("oozie.launcher." + TASK_USER_PRECEDENCE) == null) {
119            launcherJobConf.set(TASK_USER_PRECEDENCE, "true");
120        }
121        if (launcherJobConf.get("oozie.launcher." + TASK_USER_CLASSPATH_PRECEDENCE) == null) {
122            launcherJobConf.set(TASK_USER_CLASSPATH_PRECEDENCE, "true");
123        }
124        return launcherJobConf;
125    }
126
127    @Override
128    Configuration setupLauncherConf(Configuration conf, Element actionXml, Path appPath, Context context)
129            throws ActionExecutorException {
130        super.setupLauncherConf(conf, actionXml, appPath, context);
131
132        // Set SPARK_HOME environment variable on launcher job
133        // It is needed since pyspark client checks for it.
134        String sparkHome = "SPARK_HOME=.";
135        String mapredChildEnv = conf.get("oozie.launcher." + MAPRED_CHILD_ENV);
136
137        if (mapredChildEnv == null) {
138            conf.set(MAPRED_CHILD_ENV, sparkHome);
139            conf.set("oozie.launcher." + MAPRED_CHILD_ENV, sparkHome);
140        } else if (!mapredChildEnv.contains("SPARK_HOME")) {
141            conf.set(MAPRED_CHILD_ENV, mapredChildEnv + "," + sparkHome);
142            conf.set("oozie.launcher." + MAPRED_CHILD_ENV, mapredChildEnv + "," + sparkHome);
143        }
144        return conf;
145    }
146
147    @Override
148    public List<Class> getLauncherClasses() {
149        List<Class> classes = new ArrayList<Class>();
150        try {
151            classes.add(Class.forName(SPARK_MAIN_CLASS_NAME));
152        } catch (ClassNotFoundException e) {
153            throw new RuntimeException("Class not found", e);
154        }
155        return classes;
156    }
157
158
159    /**
160     * Return the sharelib name for the action.
161     *
162     * @param actionXml
163     * @return returns <code>spark</code>.
164     */
165    @Override
166    protected String getDefaultShareLibName(Element actionXml) {
167        return "spark";
168    }
169
170    @Override
171    protected String getLauncherMain(Configuration launcherConf, Element actionXml) {
172        return launcherConf.get(LauncherMapper.CONF_OOZIE_ACTION_MAIN_CLASS, SPARK_MAIN_CLASS_NAME);
173    }
174}