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.Path; 023import org.apache.hadoop.mapred.Counters; 024import org.apache.hadoop.mapred.JobClient; 025import org.apache.hadoop.mapred.JobConf; 026import org.apache.hadoop.mapred.JobID; 027import org.apache.hadoop.mapred.RunningJob; 028import org.apache.oozie.action.ActionExecutorException; 029import org.apache.oozie.client.WorkflowAction; 030import org.apache.oozie.util.XConfiguration; 031import org.apache.oozie.util.XLog; 032import org.apache.oozie.util.XmlUtils; 033import org.jdom.Element; 034import org.jdom.Namespace; 035 036import java.io.IOException; 037import java.io.StringReader; 038import java.util.ArrayList; 039import java.util.List; 040import java.util.StringTokenizer; 041 042public class SqoopActionExecutor extends JavaActionExecutor { 043 044 public static final String OOZIE_ACTION_EXTERNAL_STATS_WRITE = "oozie.action.external.stats.write"; 045 private static final String SQOOP_MAIN_CLASS_NAME = "org.apache.oozie.action.hadoop.SqoopMain"; 046 static final String SQOOP_ARGS = "oozie.sqoop.args"; 047 private static final String SQOOP = "sqoop"; 048 049 public SqoopActionExecutor() { 050 super(SQOOP); 051 } 052 053 @Override 054 public List<Class> getLauncherClasses() { 055 List<Class> classes = new ArrayList<Class>(); 056 try { 057 classes.add(Class.forName(SQOOP_MAIN_CLASS_NAME)); 058 } 059 catch (ClassNotFoundException e) { 060 throw new RuntimeException("Class not found", e); 061 } 062 return classes; 063 } 064 065 @Override 066 protected String getLauncherMain(Configuration launcherConf, Element actionXml) { 067 return launcherConf.get(LauncherMapper.CONF_OOZIE_ACTION_MAIN_CLASS, SQOOP_MAIN_CLASS_NAME); 068 } 069 070 @Override 071 @SuppressWarnings("unchecked") 072 Configuration setupActionConf(Configuration actionConf, Context context, Element actionXml, Path appPath) 073 throws ActionExecutorException { 074 super.setupActionConf(actionConf, context, actionXml, appPath); 075 Namespace ns = actionXml.getNamespace(); 076 077 try { 078 Element e = actionXml.getChild("configuration", ns); 079 if (e != null) { 080 String strConf = XmlUtils.prettyPrint(e).toString(); 081 XConfiguration inlineConf = new XConfiguration(new StringReader(strConf)); 082 checkForDisallowedProps(inlineConf, "inline configuration"); 083 XConfiguration.copy(inlineConf, actionConf); 084 } 085 } catch (IOException ex) { 086 throw convertException(ex); 087 } 088 089 final List<String> argList = new ArrayList<>(); 090 // Build a list of arguments from either a tokenized <command> string or a list of <arg> 091 if (actionXml.getChild("command", ns) != null) { 092 String command = actionXml.getChild("command", ns).getTextTrim(); 093 StringTokenizer st = new StringTokenizer(command, " "); 094 while (st.hasMoreTokens()) { 095 argList.add(st.nextToken()); 096 } 097 } 098 else { 099 List<Element> eArgs = (List<Element>) actionXml.getChildren("arg", ns); 100 for (Element elem : eArgs) { 101 argList.add(elem.getTextTrim()); 102 } 103 } 104 // If the command is given accidentally as "sqoop import --option" 105 // instead of "import --option" we can make a user's life easier 106 // by removing away the unnecessary "sqoop" token. 107 // However, we do not do this if the command looks like 108 // "sqoop --option", as that's entirely invalid. 109 if (argList.size() > 1 && 110 argList.get(0).equalsIgnoreCase(SQOOP) && 111 !argList.get(1).startsWith("-")) { 112 XLog.getLog(getClass()).info( 113 "Found a redundant 'sqoop' prefixing the command. Removing it."); 114 argList.remove(0); 115 } 116 117 setSqoopCommand(actionConf, argList.toArray(new String[argList.size()])); 118 return actionConf; 119 } 120 121 private void setSqoopCommand(Configuration conf, String[] args) { 122 MapReduceMain.setStrings(conf, SQOOP_ARGS, args); 123 } 124 125 /** 126 * We will gather counters from all executed action Hadoop jobs (e.g. jobs 127 * that moved data, not the launcher itself) and merge them together. There 128 * will be only one job most of the time. The only exception is 129 * import-all-table option that will execute one job per one exported table. 130 * 131 * @param context Action context 132 * @param action Workflow action 133 * @throws ActionExecutorException 134 */ 135 @Override 136 public void end(Context context, WorkflowAction action) throws ActionExecutorException { 137 super.end(context, action); 138 JobClient jobClient = null; 139 140 boolean exception = false; 141 try { 142 if (action.getStatus() == WorkflowAction.Status.OK) { 143 Element actionXml = XmlUtils.parseXml(action.getConf()); 144 JobConf jobConf = createBaseHadoopConf(context, actionXml); 145 jobClient = createJobClient(context, jobConf); 146 147 // Cumulative counters for all Sqoop mapreduce jobs 148 Counters counters = null; 149 150 // Sqoop do not have to create mapreduce job each time 151 String externalIds = action.getExternalChildIDs(); 152 if (externalIds != null && !externalIds.trim().isEmpty()) { 153 String []jobIds = externalIds.split(","); 154 155 for(String jobId : jobIds) { 156 RunningJob runningJob = jobClient.getJob(JobID.forName(jobId)); 157 if (runningJob == null) { 158 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "SQOOP001", 159 "Unknown hadoop job [{0}] associated with action [{1}]. Failing this action!", action 160 .getExternalId(), action.getId()); 161 } 162 163 Counters taskCounters = runningJob.getCounters(); 164 if(taskCounters != null) { 165 if(counters == null) { 166 counters = taskCounters; 167 } else { 168 counters.incrAllCounters(taskCounters); 169 } 170 } else { 171 XLog.getLog(getClass()).warn("Could not find Hadoop Counters for job: [{0}]", jobId); 172 } 173 } 174 } 175 176 if (counters != null) { 177 ActionStats stats = new MRStats(counters); 178 String statsJsonString = stats.toJSON(); 179 context.setVar(MapReduceActionExecutor.HADOOP_COUNTERS, statsJsonString); 180 181 // If action stats write property is set to false by user or 182 // size of stats is greater than the maximum allowed size, 183 // do not store the action stats 184 if (Boolean.parseBoolean(evaluateConfigurationProperty(actionXml, 185 OOZIE_ACTION_EXTERNAL_STATS_WRITE, "true")) 186 && (statsJsonString.getBytes().length <= getMaxExternalStatsSize())) { 187 context.setExecutionStats(statsJsonString); 188 LOG.debug( 189 "Printing stats for sqoop action as a JSON string : [{0}]", statsJsonString); 190 } 191 } else { 192 context.setVar(MapReduceActionExecutor.HADOOP_COUNTERS, ""); 193 XLog.getLog(getClass()).warn("Can't find any associated Hadoop job counters"); 194 } 195 } 196 } 197 catch (Exception ex) { 198 exception = true; 199 throw convertException(ex); 200 } 201 finally { 202 if (jobClient != null) { 203 try { 204 jobClient.close(); 205 } 206 catch (Exception e) { 207 if (exception) { 208 LOG.error("JobClient error: ", e); 209 } 210 else { 211 throw convertException(e); 212 } 213 } 214 } 215 } 216 } 217 218 // Return the value of the specified configuration property 219 private String evaluateConfigurationProperty(Element actionConf, String key, String defaultValue) 220 throws ActionExecutorException { 221 try { 222 if (actionConf != null) { 223 Namespace ns = actionConf.getNamespace(); 224 Element e = actionConf.getChild("configuration", ns); 225 226 if(e != null) { 227 String strConf = XmlUtils.prettyPrint(e).toString(); 228 XConfiguration inlineConf = new XConfiguration(new StringReader(strConf)); 229 return inlineConf.get(key, defaultValue); 230 } 231 } 232 return defaultValue; 233 } 234 catch (IOException ex) { 235 throw convertException(ex); 236 } 237 } 238 239 /** 240 * Return the sharelib name for the action. 241 * 242 * @return returns <code>sqoop</code>. 243 * @param actionXml 244 */ 245 @Override 246 protected String getDefaultShareLibName(Element actionXml) { 247 return "sqoop"; 248 } 249 250}