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}