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 java.io.FileNotFoundException; 022import java.io.IOException; 023import java.io.StringReader; 024import java.net.ConnectException; 025import java.net.URI; 026import java.net.URISyntaxException; 027import java.net.UnknownHostException; 028import java.security.PrivilegedExceptionAction; 029import java.text.MessageFormat; 030import java.util.ArrayList; 031import java.util.Arrays; 032import java.util.HashMap; 033import java.util.HashSet; 034import java.util.Iterator; 035import java.util.List; 036import java.util.Map; 037import java.util.Map.Entry; 038import java.util.Properties; 039import java.util.Set; 040import java.util.regex.Matcher; 041import java.util.regex.Pattern; 042 043import com.google.common.annotations.VisibleForTesting; 044import com.google.common.primitives.Ints; 045import org.apache.directory.api.util.Strings; 046import org.apache.hadoop.conf.Configuration; 047import org.apache.hadoop.filecache.DistributedCache; 048import org.apache.hadoop.fs.CommonConfigurationKeysPublic; 049import org.apache.hadoop.fs.FileStatus; 050import org.apache.hadoop.fs.FileSystem; 051import org.apache.hadoop.fs.Path; 052import org.apache.hadoop.fs.permission.AccessControlException; 053import org.apache.hadoop.mapreduce.security.token.delegation.DelegationTokenIdentifier; 054import org.apache.oozie.hadoop.utils.HadoopShims; 055import org.apache.hadoop.io.Text; 056import org.apache.hadoop.mapred.FileInputFormat; 057import org.apache.hadoop.mapred.JobClient; 058import org.apache.hadoop.mapred.JobConf; 059import org.apache.hadoop.mapred.JobID; 060import org.apache.hadoop.mapred.RunningJob; 061import org.apache.hadoop.security.UserGroupInformation; 062import org.apache.hadoop.security.token.Token; 063import org.apache.hadoop.security.token.TokenIdentifier; 064import org.apache.hadoop.util.DiskChecker; 065import org.apache.oozie.WorkflowActionBean; 066import org.apache.oozie.WorkflowJobBean; 067import org.apache.oozie.action.ActionExecutor; 068import org.apache.oozie.action.ActionExecutorException; 069import org.apache.oozie.client.OozieClient; 070import org.apache.oozie.client.WorkflowAction; 071import org.apache.oozie.client.WorkflowJob; 072import org.apache.oozie.command.coord.CoordActionStartXCommand; 073import org.apache.oozie.command.wf.WorkflowXCommand; 074import org.apache.oozie.service.ConfigurationService; 075import org.apache.oozie.service.HadoopAccessorException; 076import org.apache.oozie.service.HadoopAccessorService; 077import org.apache.oozie.service.Services; 078import org.apache.oozie.service.ShareLibService; 079import org.apache.oozie.service.URIHandlerService; 080import org.apache.oozie.service.UserGroupInformationService; 081import org.apache.oozie.service.WorkflowAppService; 082import org.apache.oozie.util.ELEvaluationException; 083import org.apache.oozie.util.ELEvaluator; 084import org.apache.oozie.util.JobUtils; 085import org.apache.oozie.util.LogUtils; 086import org.apache.oozie.util.PropertiesUtils; 087import org.apache.oozie.util.XConfiguration; 088import org.apache.oozie.util.XLog; 089import org.apache.oozie.util.XmlUtils; 090import org.jdom.Element; 091import org.jdom.JDOMException; 092import org.jdom.Namespace; 093 094 095public class JavaActionExecutor extends ActionExecutor { 096 097 private static final String HADOOP_USER = "user.name"; 098 public static final String HADOOP_JOB_TRACKER = "mapred.job.tracker"; 099 public static final String HADOOP_JOB_TRACKER_2 = "mapreduce.jobtracker.address"; 100 public static final String HADOOP_YARN_RM = "yarn.resourcemanager.address"; 101 public static final String HADOOP_NAME_NODE = "fs.default.name"; 102 private static final String HADOOP_JOB_NAME = "mapred.job.name"; 103 public static final String OOZIE_COMMON_LIBDIR = "oozie"; 104 private static final Set<String> DISALLOWED_PROPERTIES = new HashSet<String>(); 105 public final static String MAX_EXTERNAL_STATS_SIZE = "oozie.external.stats.max.size"; 106 public static final String ACL_VIEW_JOB = "mapreduce.job.acl-view-job"; 107 public static final String ACL_MODIFY_JOB = "mapreduce.job.acl-modify-job"; 108 public static final String HADOOP_YARN_UBER_MODE = "mapreduce.job.ubertask.enable"; 109 public static final String OOZIE_ACTION_LAUNCHER_PREFIX = ActionExecutor.CONF_PREFIX + "launcher."; 110 public static final String HADOOP_YARN_KILL_CHILD_JOBS_ON_AMRESTART = 111 OOZIE_ACTION_LAUNCHER_PREFIX + "am.restart.kill.childjobs"; 112 public static final String HADOOP_MAP_MEMORY_MB = "mapreduce.map.memory.mb"; 113 public static final String HADOOP_CHILD_JAVA_OPTS = "mapred.child.java.opts"; 114 public static final String HADOOP_MAP_JAVA_OPTS = "mapreduce.map.java.opts"; 115 public static final String HADOOP_REDUCE_JAVA_OPTS = "mapreduce.reduce.java.opts"; 116 public static final String HADOOP_CHILD_JAVA_ENV = "mapred.child.env"; 117 public static final String HADOOP_MAP_JAVA_ENV = "mapreduce.map.env"; 118 public static final String YARN_AM_RESOURCE_MB = "yarn.app.mapreduce.am.resource.mb"; 119 public static final String YARN_AM_COMMAND_OPTS = "yarn.app.mapreduce.am.command-opts"; 120 public static final String YARN_AM_ENV = "yarn.app.mapreduce.am.env"; 121 private static final String JAVA_MAIN_CLASS_NAME = "org.apache.oozie.action.hadoop.JavaMain"; 122 public static final int YARN_MEMORY_MB_MIN = 512; 123 private static int maxActionOutputLen; 124 private static int maxExternalStatsSize; 125 private static int maxFSGlobMax; 126 private static final String SUCCEEDED = "SUCCEEDED"; 127 private static final String KILLED = "KILLED"; 128 private static final String FAILED = "FAILED"; 129 private static final String FAILED_KILLED = "FAILED/KILLED"; 130 protected XLog LOG = XLog.getLog(getClass()); 131 private static final Pattern heapPattern = Pattern.compile("-Xmx(([0-9]+)[mMgG])"); 132 private static final String JAVA_TMP_DIR_SETTINGS = "-Djava.io.tmpdir="; 133 public static final String CONF_HADOOP_YARN_UBER_MODE = OOZIE_ACTION_LAUNCHER_PREFIX + HADOOP_YARN_UBER_MODE; 134 public static final String HADOOP_JOB_CLASSLOADER = "mapreduce.job.classloader"; 135 public static final String HADOOP_USER_CLASSPATH_FIRST = "mapreduce.user.classpath.first"; 136 public static final String OOZIE_CREDENTIALS_SKIP = "oozie.credentials.skip"; 137 private static final LauncherInputFormatClassLocator launcherInputFormatClassLocator = new LauncherInputFormatClassLocator(); 138 139 public XConfiguration workflowConf = null; 140 141 static { 142 DISALLOWED_PROPERTIES.add(HADOOP_USER); 143 DISALLOWED_PROPERTIES.add(HADOOP_JOB_TRACKER); 144 DISALLOWED_PROPERTIES.add(HADOOP_NAME_NODE); 145 DISALLOWED_PROPERTIES.add(HADOOP_JOB_TRACKER_2); 146 DISALLOWED_PROPERTIES.add(HADOOP_YARN_RM); 147 } 148 149 public JavaActionExecutor() { 150 this("java"); 151 } 152 153 protected JavaActionExecutor(String type) { 154 super(type); 155 } 156 157 public static List<Class> getCommonLauncherClasses() { 158 List<Class> classes = new ArrayList<Class>(); 159 classes.add(LauncherMapper.class); 160 classes.add(launcherInputFormatClassLocator.locateOrGet()); 161 classes.add(OozieLauncherOutputFormat.class); 162 classes.add(OozieLauncherOutputCommitter.class); 163 classes.add(LauncherMainHadoopUtils.class); 164 classes.add(HadoopShims.class); 165 classes.addAll(Services.get().get(URIHandlerService.class).getClassesForLauncher()); 166 return classes; 167 } 168 169 public List<Class> getLauncherClasses() { 170 List<Class> classes = new ArrayList<Class>(); 171 try { 172 classes.add(Class.forName(JAVA_MAIN_CLASS_NAME)); 173 } 174 catch (ClassNotFoundException e) { 175 throw new RuntimeException("Class not found", e); 176 } 177 return classes; 178 } 179 180 @Override 181 public void initActionType() { 182 super.initActionType(); 183 maxActionOutputLen = ConfigurationService.getInt(LauncherMapper.CONF_OOZIE_ACTION_MAX_OUTPUT_DATA); 184 //Get the limit for the maximum allowed size of action stats 185 maxExternalStatsSize = ConfigurationService.getInt(JavaActionExecutor.MAX_EXTERNAL_STATS_SIZE); 186 maxExternalStatsSize = (maxExternalStatsSize == -1) ? Integer.MAX_VALUE : maxExternalStatsSize; 187 //Get the limit for the maximum number of globbed files/dirs for FS operation 188 maxFSGlobMax = ConfigurationService.getInt(LauncherMapper.CONF_OOZIE_ACTION_FS_GLOB_MAX); 189 190 registerError(UnknownHostException.class.getName(), ActionExecutorException.ErrorType.TRANSIENT, "JA001"); 191 registerError(AccessControlException.class.getName(), ActionExecutorException.ErrorType.NON_TRANSIENT, 192 "JA002"); 193 registerError(DiskChecker.DiskOutOfSpaceException.class.getName(), 194 ActionExecutorException.ErrorType.NON_TRANSIENT, "JA003"); 195 registerError(org.apache.hadoop.hdfs.protocol.QuotaExceededException.class.getName(), 196 ActionExecutorException.ErrorType.NON_TRANSIENT, "JA004"); 197 registerError(org.apache.hadoop.hdfs.server.namenode.SafeModeException.class.getName(), 198 ActionExecutorException.ErrorType.NON_TRANSIENT, "JA005"); 199 registerError(ConnectException.class.getName(), ActionExecutorException.ErrorType.TRANSIENT, " JA006"); 200 registerError(JDOMException.class.getName(), ActionExecutorException.ErrorType.ERROR, "JA007"); 201 registerError(FileNotFoundException.class.getName(), ActionExecutorException.ErrorType.ERROR, "JA008"); 202 registerError(IOException.class.getName(), ActionExecutorException.ErrorType.TRANSIENT, "JA009"); 203 } 204 205 206 /** 207 * Get the maximum allowed size of stats 208 * 209 * @return maximum size of stats 210 */ 211 public static int getMaxExternalStatsSize() { 212 return maxExternalStatsSize; 213 } 214 215 static void checkForDisallowedProps(Configuration conf, String confName) throws ActionExecutorException { 216 for (String prop : DISALLOWED_PROPERTIES) { 217 if (conf.get(prop) != null) { 218 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "JA010", 219 "Property [{0}] not allowed in action [{1}] configuration", prop, confName); 220 } 221 } 222 } 223 224 public JobConf createBaseHadoopConf(Context context, Element actionXml) { 225 Namespace ns = actionXml.getNamespace(); 226 String jobTracker = actionXml.getChild("job-tracker", ns).getTextTrim(); 227 String nameNode = actionXml.getChild("name-node", ns).getTextTrim(); 228 JobConf conf = Services.get().get(HadoopAccessorService.class).createJobConf(jobTracker); 229 conf.set(HADOOP_USER, context.getProtoActionConf().get(WorkflowAppService.HADOOP_USER)); 230 conf.set(HADOOP_JOB_TRACKER, jobTracker); 231 conf.set(HADOOP_JOB_TRACKER_2, jobTracker); 232 conf.set(HADOOP_YARN_RM, jobTracker); 233 conf.set(HADOOP_NAME_NODE, nameNode); 234 conf.set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "true"); 235 return conf; 236 } 237 238 private static void injectLauncherProperties(Configuration srcConf, Configuration launcherConf) { 239 for (Map.Entry<String, String> entry : srcConf) { 240 if (entry.getKey().startsWith("oozie.launcher.")) { 241 String name = entry.getKey().substring("oozie.launcher.".length()); 242 String value = entry.getValue(); 243 // setting original KEY 244 launcherConf.set(entry.getKey(), value); 245 // setting un-prefixed key (to allow Hadoop job config 246 // for the launcher job 247 launcherConf.set(name, value); 248 } 249 } 250 } 251 252 Configuration setupLauncherConf(Configuration conf, Element actionXml, Path appPath, Context context) 253 throws ActionExecutorException { 254 try { 255 Namespace ns = actionXml.getNamespace(); 256 XConfiguration launcherConf = new XConfiguration(); 257 // Inject action defaults for launcher 258 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 259 XConfiguration actionDefaultConf = has.createActionDefaultConf(conf.get(HADOOP_JOB_TRACKER), getType()); 260 injectLauncherProperties(actionDefaultConf, launcherConf); 261 // Inject <job-xml> and <configuration> for launcher 262 try { 263 parseJobXmlAndConfiguration(context, actionXml, appPath, conf, true); 264 } catch (HadoopAccessorException ex) { 265 throw convertException(ex); 266 } catch (URISyntaxException ex) { 267 throw convertException(ex); 268 } 269 // Inject use uber mode for launcher 270 injectLauncherUseUberMode(launcherConf); 271 XConfiguration.copy(launcherConf, conf); 272 checkForDisallowedProps(launcherConf, "launcher configuration"); 273 // Inject config-class for launcher to use for action 274 Element e = actionXml.getChild("config-class", ns); 275 if (e != null) { 276 conf.set(LauncherMapper.OOZIE_ACTION_CONFIG_CLASS, e.getTextTrim()); 277 } 278 return conf; 279 } 280 catch (IOException ex) { 281 throw convertException(ex); 282 } 283 } 284 285 void injectLauncherUseUberMode(Configuration launcherConf) { 286 // Set Uber Mode for the launcher (YARN only, ignored by MR1) 287 // Priority: 288 // 1. action's <configuration> 289 // 2. oozie.action.#action-type#.launcher.mapreduce.job.ubertask.enable 290 // 3. oozie.action.launcher.mapreduce.job.ubertask.enable 291 if (launcherConf.get(HADOOP_YARN_UBER_MODE) == null) { 292 if (ConfigurationService.get(getActionTypeLauncherPrefix() + HADOOP_YARN_UBER_MODE).length() > 0) { 293 if (ConfigurationService.getBoolean(getActionTypeLauncherPrefix() + HADOOP_YARN_UBER_MODE)) { 294 launcherConf.setBoolean(HADOOP_YARN_UBER_MODE, true); 295 } 296 } else { 297 if (ConfigurationService.getBoolean(OOZIE_ACTION_LAUNCHER_PREFIX + HADOOP_YARN_UBER_MODE)) { 298 launcherConf.setBoolean(HADOOP_YARN_UBER_MODE, true); 299 } 300 } 301 } 302 } 303 304 void updateConfForUberMode(Configuration launcherConf) { 305 306 // child.env 307 boolean hasConflictEnv = false; 308 String launcherMapEnv = launcherConf.get(HADOOP_MAP_JAVA_ENV); 309 if (launcherMapEnv == null) { 310 launcherMapEnv = launcherConf.get(HADOOP_CHILD_JAVA_ENV); 311 } 312 String amEnv = launcherConf.get(YARN_AM_ENV); 313 StringBuffer envStr = new StringBuffer(); 314 HashMap<String, List<String>> amEnvMap = null; 315 HashMap<String, List<String>> launcherMapEnvMap = null; 316 if (amEnv != null) { 317 envStr.append(amEnv); 318 amEnvMap = populateEnvMap(amEnv); 319 } 320 if (launcherMapEnv != null) { 321 launcherMapEnvMap = populateEnvMap(launcherMapEnv); 322 if (amEnvMap != null) { 323 Iterator<String> envKeyItr = launcherMapEnvMap.keySet().iterator(); 324 while (envKeyItr.hasNext()) { 325 String envKey = envKeyItr.next(); 326 if (amEnvMap.containsKey(envKey)) { 327 List<String> amValList = amEnvMap.get(envKey); 328 List<String> launcherValList = launcherMapEnvMap.get(envKey); 329 Iterator<String> valItr = launcherValList.iterator(); 330 while (valItr.hasNext()) { 331 String val = valItr.next(); 332 if (!amValList.contains(val)) { 333 hasConflictEnv = true; 334 break; 335 } 336 else { 337 valItr.remove(); 338 } 339 } 340 if (launcherValList.isEmpty()) { 341 envKeyItr.remove(); 342 } 343 } 344 } 345 } 346 } 347 if (hasConflictEnv) { 348 launcherConf.setBoolean(HADOOP_YARN_UBER_MODE, false); 349 } 350 else { 351 if (launcherMapEnvMap != null) { 352 for (String key : launcherMapEnvMap.keySet()) { 353 List<String> launcherValList = launcherMapEnvMap.get(key); 354 for (String val : launcherValList) { 355 if (envStr.length() > 0) { 356 envStr.append(","); 357 } 358 envStr.append(key).append("=").append(val); 359 } 360 } 361 } 362 363 launcherConf.set(YARN_AM_ENV, envStr.toString()); 364 365 // memory.mb 366 int launcherMapMemoryMB = launcherConf.getInt(HADOOP_MAP_MEMORY_MB, 1536); 367 int amMemoryMB = launcherConf.getInt(YARN_AM_RESOURCE_MB, 1536); 368 // YARN_MEMORY_MB_MIN to provide buffer. 369 // suppose launcher map aggressively use high memory, need some 370 // headroom for AM 371 int memoryMB = Math.max(launcherMapMemoryMB, amMemoryMB) + YARN_MEMORY_MB_MIN; 372 // limit to 4096 in case of 32 bit 373 if (launcherMapMemoryMB < 4096 && amMemoryMB < 4096 && memoryMB > 4096) { 374 memoryMB = 4096; 375 } 376 launcherConf.setInt(YARN_AM_RESOURCE_MB, memoryMB); 377 378 // We already made mapred.child.java.opts and 379 // mapreduce.map.java.opts equal, so just start with one of them 380 String launcherMapOpts = launcherConf.get(HADOOP_MAP_JAVA_OPTS, ""); 381 String amChildOpts = launcherConf.get(YARN_AM_COMMAND_OPTS); 382 StringBuilder optsStr = new StringBuilder(); 383 int heapSizeForMap = extractHeapSizeMB(launcherMapOpts); 384 int heapSizeForAm = extractHeapSizeMB(amChildOpts); 385 int heapSize = Math.max(heapSizeForMap, heapSizeForAm) + YARN_MEMORY_MB_MIN; 386 // limit to 3584 in case of 32 bit 387 if (heapSizeForMap < 4096 && heapSizeForAm < 4096 && heapSize > 3584) { 388 heapSize = 3584; 389 } 390 if (amChildOpts != null) { 391 optsStr.append(amChildOpts); 392 } 393 optsStr.append(" ").append(launcherMapOpts.trim()); 394 if (heapSize > 0) { 395 // append calculated total heap size to the end 396 optsStr.append(" ").append("-Xmx").append(heapSize).append("m"); 397 } 398 launcherConf.set(YARN_AM_COMMAND_OPTS, optsStr.toString().trim()); 399 } 400 } 401 402 void updateConfForJavaTmpDir(Configuration conf) { 403 String amChildOpts = conf.get(YARN_AM_COMMAND_OPTS); 404 String oozieJavaTmpDirSetting = "-Djava.io.tmpdir=./tmp"; 405 if (amChildOpts != null && !amChildOpts.contains(JAVA_TMP_DIR_SETTINGS)) { 406 conf.set(YARN_AM_COMMAND_OPTS, amChildOpts + " " + oozieJavaTmpDirSetting); 407 } 408 } 409 410 private HashMap<String, List<String>> populateEnvMap(String input) { 411 HashMap<String, List<String>> envMaps = new HashMap<String, List<String>>(); 412 String[] envEntries = input.split(","); 413 for (String envEntry : envEntries) { 414 String[] envKeyVal = envEntry.split("="); 415 String envKey = envKeyVal[0].trim(); 416 List<String> valList = envMaps.get(envKey); 417 if (valList == null) { 418 valList = new ArrayList<String>(); 419 } 420 valList.add(envKeyVal[1].trim()); 421 envMaps.put(envKey, valList); 422 } 423 return envMaps; 424 } 425 426 public int extractHeapSizeMB(String input) { 427 int ret = 0; 428 if(input == null || input.equals("")) 429 return ret; 430 Matcher m = heapPattern.matcher(input); 431 String heapStr = null; 432 String heapNum = null; 433 // Grabs the last match which takes effect (in case that multiple Xmx options specified) 434 while (m.find()) { 435 heapStr = m.group(1); 436 heapNum = m.group(2); 437 } 438 if (heapStr != null) { 439 // when Xmx specified in Gigabyte 440 if(heapStr.endsWith("g") || heapStr.endsWith("G")) { 441 ret = Integer.parseInt(heapNum) * 1024; 442 } else { 443 ret = Integer.parseInt(heapNum); 444 } 445 } 446 return ret; 447 } 448 449 public static void parseJobXmlAndConfiguration(Context context, Element element, Path appPath, Configuration conf) 450 throws IOException, ActionExecutorException, HadoopAccessorException, URISyntaxException { 451 parseJobXmlAndConfiguration(context, element, appPath, conf, false); 452 } 453 454 public static void parseJobXmlAndConfiguration(Context context, Element element, Path appPath, Configuration conf, 455 boolean isLauncher) throws IOException, ActionExecutorException, HadoopAccessorException, URISyntaxException { 456 Namespace ns = element.getNamespace(); 457 Iterator<Element> it = element.getChildren("job-xml", ns).iterator(); 458 HashMap<String, FileSystem> filesystemsMap = new HashMap<String, FileSystem>(); 459 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 460 while (it.hasNext()) { 461 Element e = it.next(); 462 String jobXml = e.getTextTrim(); 463 Path pathSpecified = new Path(jobXml); 464 Path path = pathSpecified.isAbsolute() ? pathSpecified : new Path(appPath, jobXml); 465 FileSystem fs; 466 if (filesystemsMap.containsKey(path.toUri().getAuthority())) { 467 fs = filesystemsMap.get(path.toUri().getAuthority()); 468 } 469 else { 470 if (path.toUri().getAuthority() != null) { 471 fs = has.createFileSystem(context.getWorkflow().getUser(), path.toUri(), 472 has.createJobConf(path.toUri().getAuthority())); 473 } 474 else { 475 fs = context.getAppFileSystem(); 476 } 477 filesystemsMap.put(path.toUri().getAuthority(), fs); 478 } 479 Configuration jobXmlConf = new XConfiguration(fs.open(path)); 480 try { 481 String jobXmlConfString = XmlUtils.prettyPrint(jobXmlConf).toString(); 482 jobXmlConfString = XmlUtils.removeComments(jobXmlConfString); 483 jobXmlConfString = context.getELEvaluator().evaluate(jobXmlConfString, String.class); 484 jobXmlConf = new XConfiguration(new StringReader(jobXmlConfString)); 485 } 486 catch (ELEvaluationException ex) { 487 throw new ActionExecutorException(ActionExecutorException.ErrorType.TRANSIENT, "EL_EVAL_ERROR", ex 488 .getMessage(), ex); 489 } 490 catch (Exception ex) { 491 context.setErrorInfo("EL_ERROR", ex.getMessage()); 492 } 493 checkForDisallowedProps(jobXmlConf, "job-xml"); 494 if (isLauncher) { 495 injectLauncherProperties(jobXmlConf, conf); 496 } else { 497 XConfiguration.copy(jobXmlConf, conf); 498 } 499 } 500 Element e = element.getChild("configuration", ns); 501 if (e != null) { 502 String strConf = XmlUtils.prettyPrint(e).toString(); 503 XConfiguration inlineConf = new XConfiguration(new StringReader(strConf)); 504 checkForDisallowedProps(inlineConf, "inline configuration"); 505 if (isLauncher) { 506 injectLauncherProperties(inlineConf, conf); 507 } else { 508 XConfiguration.copy(inlineConf, conf); 509 } 510 } 511 } 512 513 Configuration setupActionConf(Configuration actionConf, Context context, Element actionXml, Path appPath) 514 throws ActionExecutorException { 515 try { 516 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 517 XConfiguration actionDefaults = has.createActionDefaultConf(actionConf.get(HADOOP_JOB_TRACKER), getType()); 518 XConfiguration.injectDefaults(actionDefaults, actionConf); 519 520 has.checkSupportedFilesystem(appPath.toUri()); 521 522 // Set the Java Main Class for the Java action to give to the Java launcher 523 setJavaMain(actionConf, actionXml); 524 525 parseJobXmlAndConfiguration(context, actionXml, appPath, actionConf); 526 527 // set cancel.delegation.token in actionConf that child job doesn't cancel delegation token 528 actionConf.setBoolean("mapreduce.job.complete.cancel.delegation.tokens", false); 529 updateConfForJavaTmpDir(actionConf); 530 setRootLoggerLevel(actionConf); 531 return actionConf; 532 } 533 catch (IOException ex) { 534 throw convertException(ex); 535 } 536 catch (HadoopAccessorException ex) { 537 throw convertException(ex); 538 } 539 catch (URISyntaxException ex) { 540 throw convertException(ex); 541 } 542 } 543 544 /** 545 * Set root log level property in actionConf 546 * @param actionConf 547 */ 548 void setRootLoggerLevel(Configuration actionConf) { 549 String oozieActionTypeRootLogger = "oozie.action." + getType() + LauncherMapper.ROOT_LOGGER_LEVEL; 550 String oozieActionRootLogger = "oozie.action." + LauncherMapper.ROOT_LOGGER_LEVEL; 551 552 // check if root log level has already mentioned in action configuration 553 String rootLogLevel = actionConf.get(oozieActionTypeRootLogger, actionConf.get(oozieActionRootLogger)); 554 if (rootLogLevel != null) { 555 // root log level is mentioned in action configuration 556 return; 557 } 558 559 // set the root log level which is mentioned in oozie default 560 rootLogLevel = ConfigurationService.get(oozieActionTypeRootLogger); 561 if (rootLogLevel != null && rootLogLevel.length() > 0) { 562 actionConf.set(oozieActionRootLogger, rootLogLevel); 563 } 564 else { 565 rootLogLevel = ConfigurationService.get(oozieActionRootLogger); 566 if (rootLogLevel != null && rootLogLevel.length() > 0) { 567 actionConf.set(oozieActionRootLogger, rootLogLevel); 568 } 569 } 570 } 571 572 Configuration addToCache(Configuration conf, Path appPath, String filePath, boolean archive) 573 throws ActionExecutorException { 574 575 URI uri = null; 576 try { 577 uri = new URI(getTrimmedEncodedPath(filePath)); 578 URI baseUri = appPath.toUri(); 579 if (uri.getScheme() == null) { 580 String resolvedPath = uri.getPath(); 581 if (!resolvedPath.startsWith("/")) { 582 resolvedPath = baseUri.getPath() + "/" + resolvedPath; 583 } 584 uri = new URI(baseUri.getScheme(), baseUri.getAuthority(), resolvedPath, uri.getQuery(), uri.getFragment()); 585 } 586 if (archive) { 587 DistributedCache.addCacheArchive(uri.normalize(), conf); 588 } 589 else { 590 String fileName = filePath.substring(filePath.lastIndexOf("/") + 1); 591 if (fileName.endsWith(".so") || fileName.contains(".so.")) { // .so files 592 uri = new URI(uri.getScheme(), uri.getAuthority(), uri.getPath(), uri.getQuery(), fileName); 593 DistributedCache.addCacheFile(uri.normalize(), conf); 594 } 595 else if (fileName.endsWith(".jar")) { // .jar files 596 if (!fileName.contains("#")) { 597 String user = conf.get("user.name"); 598 Path pathToAdd = new Path(uri.normalize()); 599 Services.get().get(HadoopAccessorService.class).addFileToClassPath(user, pathToAdd, conf); 600 } 601 else { 602 DistributedCache.addCacheFile(uri.normalize(), conf); 603 } 604 } 605 else { // regular files 606 if (!fileName.contains("#")) { 607 uri = new URI(uri.getScheme(), uri.getAuthority(), uri.getPath(), uri.getQuery(), fileName); 608 } 609 DistributedCache.addCacheFile(uri.normalize(), conf); 610 } 611 } 612 DistributedCache.createSymlink(conf); 613 return conf; 614 } 615 catch (Exception ex) { 616 LOG.debug( 617 "Errors when add to DistributedCache. Path=" + uri.toString() + ", archive=" + archive + ", conf=" 618 + XmlUtils.prettyPrint(conf).toString()); 619 throw convertException(ex); 620 } 621 } 622 623 public void prepareActionDir(FileSystem actionFs, Context context) throws ActionExecutorException { 624 try { 625 Path actionDir = context.getActionDir(); 626 Path tempActionDir = new Path(actionDir.getParent(), actionDir.getName() + ".tmp"); 627 if (!actionFs.exists(actionDir)) { 628 try { 629 actionFs.mkdirs(tempActionDir); 630 actionFs.rename(tempActionDir, actionDir); 631 } 632 catch (IOException ex) { 633 actionFs.delete(tempActionDir, true); 634 actionFs.delete(actionDir, true); 635 throw ex; 636 } 637 } 638 } 639 catch (Exception ex) { 640 throw convertException(ex); 641 } 642 } 643 644 void cleanUpActionDir(FileSystem actionFs, Context context) throws ActionExecutorException { 645 try { 646 Path actionDir = context.getActionDir(); 647 if (!context.getProtoActionConf().getBoolean(WorkflowXCommand.KEEP_WF_ACTION_DIR, false) 648 && actionFs.exists(actionDir)) { 649 actionFs.delete(actionDir, true); 650 } 651 } 652 catch (Exception ex) { 653 throw convertException(ex); 654 } 655 } 656 657 protected void addShareLib(Configuration conf, String[] actionShareLibNames) 658 throws ActionExecutorException { 659 Set<String> confSet = new HashSet<String>(Arrays.asList(getShareLibFilesForActionConf() == null ? new String[0] 660 : getShareLibFilesForActionConf())); 661 662 Set<Path> sharelibList = new HashSet<Path>(); 663 664 if (actionShareLibNames != null) { 665 try { 666 ShareLibService shareLibService = Services.get().get(ShareLibService.class); 667 FileSystem fs = shareLibService.getFileSystem(); 668 if (fs != null) { 669 for (String actionShareLibName : actionShareLibNames) { 670 List<Path> listOfPaths = shareLibService.getShareLibJars(actionShareLibName); 671 if (listOfPaths != null && !listOfPaths.isEmpty()) { 672 for (Path actionLibPath : listOfPaths) { 673 String fragmentName = new URI(actionLibPath.toString()).getFragment(); 674 Path pathWithFragment = fragmentName == null ? actionLibPath : new Path(new URI( 675 actionLibPath.toString()).getPath()); 676 String fileName = fragmentName == null ? actionLibPath.getName() : fragmentName; 677 if (confSet.contains(fileName)) { 678 Configuration jobXmlConf = shareLibService.getShareLibConf(actionShareLibName, 679 pathWithFragment); 680 if (jobXmlConf != null) { 681 checkForDisallowedProps(jobXmlConf, actionLibPath.getName()); 682 XConfiguration.injectDefaults(jobXmlConf, conf); 683 LOG.trace("Adding properties of " + actionLibPath + " to job conf"); 684 } 685 } 686 else { 687 // Filtering out duplicate jars or files 688 sharelibList.add(new Path(actionLibPath.toUri()) { 689 @Override 690 public int hashCode() { 691 return getName().hashCode(); 692 } 693 @Override 694 public String getName() { 695 try { 696 return (new URI(toString())).getFragment() == null ? new Path(toUri()).getName() 697 : (new URI(toString())).getFragment(); 698 } 699 catch (URISyntaxException e) { 700 throw new RuntimeException(e); 701 } 702 } 703 @Override 704 public boolean equals(Object input) { 705 if (input == null) { 706 return false; 707 } 708 if (input == this) { 709 return true; 710 } 711 if (!(input instanceof Path)) { 712 return false; 713 } 714 return getName().equals(((Path) input).getName()); 715 } 716 }); 717 } 718 } 719 } 720 } 721 } 722 for (Path libPath : sharelibList) { 723 addToCache(conf, libPath, libPath.toUri().getPath(), false); 724 } 725 } 726 catch (URISyntaxException ex) { 727 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "Error configuring sharelib", 728 ex.getMessage()); 729 } 730 catch (IOException ex) { 731 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "It should never happen", 732 ex.getMessage()); 733 } 734 } 735 } 736 737 protected void addSystemShareLibForAction(Configuration conf) throws ActionExecutorException { 738 ShareLibService shareLibService = Services.get().get(ShareLibService.class); 739 // ShareLibService is null for test cases 740 if (shareLibService != null) { 741 try { 742 List<Path> listOfPaths = shareLibService.getSystemLibJars(JavaActionExecutor.OOZIE_COMMON_LIBDIR); 743 if (listOfPaths == null || listOfPaths.isEmpty()) { 744 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "EJ001", 745 "Could not locate Oozie sharelib"); 746 } 747 FileSystem fs = listOfPaths.get(0).getFileSystem(conf); 748 for (Path actionLibPath : listOfPaths) { 749 JobUtils.addFileToClassPath(actionLibPath, conf, fs); 750 DistributedCache.createSymlink(conf); 751 } 752 listOfPaths = shareLibService.getSystemLibJars(getType()); 753 if (listOfPaths != null) { 754 for (Path actionLibPath : listOfPaths) { 755 JobUtils.addFileToClassPath(actionLibPath, conf, fs); 756 DistributedCache.createSymlink(conf); 757 } 758 } 759 } 760 catch (IOException ex) { 761 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "It should never happen", 762 ex.getMessage()); 763 } 764 } 765 } 766 767 protected void addActionLibs(Path appPath, Configuration conf) throws ActionExecutorException { 768 String[] actionLibsStrArr = conf.getStrings("oozie.launcher.oozie.libpath"); 769 if (actionLibsStrArr != null) { 770 try { 771 for (String actionLibsStr : actionLibsStrArr) { 772 actionLibsStr = actionLibsStr.trim(); 773 if (actionLibsStr.length() > 0) 774 { 775 Path actionLibsPath = new Path(actionLibsStr); 776 String user = conf.get("user.name"); 777 FileSystem fs = Services.get().get(HadoopAccessorService.class).createFileSystem(user, appPath.toUri(), conf); 778 if (fs.exists(actionLibsPath)) { 779 FileStatus[] files = fs.listStatus(actionLibsPath); 780 for (FileStatus file : files) { 781 addToCache(conf, appPath, file.getPath().toUri().getPath(), false); 782 } 783 } 784 } 785 } 786 } 787 catch (HadoopAccessorException ex){ 788 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, 789 ex.getErrorCode().toString(), ex.getMessage()); 790 } 791 catch (IOException ex){ 792 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, 793 "It should never happen", ex.getMessage()); 794 } 795 } 796 } 797 798 @SuppressWarnings("unchecked") 799 public void setLibFilesArchives(Context context, Element actionXml, Path appPath, Configuration conf) 800 throws ActionExecutorException { 801 Configuration proto = context.getProtoActionConf(); 802 803 // Workflow lib/ 804 String[] paths = proto.getStrings(WorkflowAppService.APP_LIB_PATH_LIST); 805 if (paths != null) { 806 for (String path : paths) { 807 addToCache(conf, appPath, path, false); 808 } 809 } 810 811 // Action libs 812 addActionLibs(appPath, conf); 813 814 // files and archives defined in the action 815 for (Element eProp : (List<Element>) actionXml.getChildren()) { 816 if (eProp.getName().equals("file")) { 817 String[] filePaths = eProp.getTextTrim().split(","); 818 for (String path : filePaths) { 819 addToCache(conf, appPath, path, false); 820 } 821 } 822 else if (eProp.getName().equals("archive")) { 823 String[] archivePaths = eProp.getTextTrim().split(","); 824 for (String path : archivePaths){ 825 addToCache(conf, appPath, path.trim(), true); 826 } 827 } 828 } 829 830 addAllShareLibs(appPath, conf, context, actionXml); 831 } 832 833 @VisibleForTesting 834 protected static String getTrimmedEncodedPath(String path) { 835 return path.trim().replace(" ", "%20"); 836 } 837 838 // Adds action specific share libs and common share libs 839 private void addAllShareLibs(Path appPath, Configuration conf, Context context, Element actionXml) 840 throws ActionExecutorException { 841 // Add action specific share libs 842 addActionShareLib(appPath, conf, context, actionXml); 843 // Add common sharelibs for Oozie and launcher jars 844 addSystemShareLibForAction(conf); 845 } 846 847 private void addActionShareLib(Path appPath, Configuration conf, Context context, Element actionXml) 848 throws ActionExecutorException { 849 XConfiguration wfJobConf = null; 850 try { 851 wfJobConf = getWorkflowConf(context); 852 } 853 catch (IOException ioe) { 854 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "It should never happen", 855 ioe.getMessage()); 856 } 857 // Action sharelibs are only added if user has specified to use system libpath 858 if (conf.get(OozieClient.USE_SYSTEM_LIBPATH) == null) { 859 if (wfJobConf.getBoolean(OozieClient.USE_SYSTEM_LIBPATH, 860 ConfigurationService.getBoolean(OozieClient.USE_SYSTEM_LIBPATH))) { 861 // add action specific sharelibs 862 addShareLib(conf, getShareLibNames(context, actionXml, conf)); 863 } 864 } 865 else { 866 if (conf.getBoolean(OozieClient.USE_SYSTEM_LIBPATH, false)) { 867 // add action specific sharelibs 868 addShareLib(conf, getShareLibNames(context, actionXml, conf)); 869 } 870 } 871 } 872 873 874 protected String getLauncherMain(Configuration launcherConf, Element actionXml) { 875 return launcherConf.get(LauncherMapper.CONF_OOZIE_ACTION_MAIN_CLASS, JavaMain.class.getName()); 876 } 877 878 private void setJavaMain(Configuration actionConf, Element actionXml) { 879 Namespace ns = actionXml.getNamespace(); 880 Element e = actionXml.getChild("main-class", ns); 881 if (e != null) { 882 actionConf.set(JavaMain.JAVA_MAIN_CLASS, e.getTextTrim()); 883 } 884 } 885 886 private static final String QUEUE_NAME = "mapred.job.queue.name"; 887 888 private static final Set<String> SPECIAL_PROPERTIES = new HashSet<String>(); 889 890 static { 891 SPECIAL_PROPERTIES.add(QUEUE_NAME); 892 SPECIAL_PROPERTIES.add(ACL_VIEW_JOB); 893 SPECIAL_PROPERTIES.add(ACL_MODIFY_JOB); 894 } 895 896 @SuppressWarnings("unchecked") 897 JobConf createLauncherConf(FileSystem actionFs, Context context, WorkflowAction action, Element actionXml, Configuration actionConf) 898 throws ActionExecutorException { 899 try { 900 901 // app path could be a file 902 Path appPathRoot = new Path(context.getWorkflow().getAppPath()); 903 if (actionFs.isFile(appPathRoot)) { 904 appPathRoot = appPathRoot.getParent(); 905 } 906 907 // launcher job configuration 908 JobConf launcherJobConf = createBaseHadoopConf(context, actionXml); 909 // cancel delegation token on a launcher job which stays alive till child job(s) finishes 910 // otherwise (in mapred action), doesn't cancel not to disturb running child job 911 launcherJobConf.setBoolean("mapreduce.job.complete.cancel.delegation.tokens", true); 912 setupLauncherConf(launcherJobConf, actionXml, appPathRoot, context); 913 914 // Properties for when a launcher job's AM gets restarted 915 if (ConfigurationService.getBoolean(HADOOP_YARN_KILL_CHILD_JOBS_ON_AMRESTART)) { 916 // launcher time filter is required to prune the search of launcher tag. 917 // Setting coordinator action nominal time as launcher time as it child job cannot launch before nominal 918 // time. Workflow created time is good enough when workflow is running independently or workflow is 919 // rerunning from failed node. 920 long launcherTime = System.currentTimeMillis(); 921 String coordActionNominalTime = context.getProtoActionConf().get( 922 CoordActionStartXCommand.OOZIE_COORD_ACTION_NOMINAL_TIME); 923 if (coordActionNominalTime != null) { 924 launcherTime = Long.parseLong(coordActionNominalTime); 925 } 926 else if (context.getWorkflow().getCreatedTime() != null) { 927 launcherTime = context.getWorkflow().getCreatedTime().getTime(); 928 } 929 String actionYarnTag = getActionYarnTag(getWorkflowConf(context), context.getWorkflow(), action); 930 LauncherMapperHelper.setupYarnRestartHandling(launcherJobConf, actionConf, actionYarnTag, launcherTime); 931 } 932 else { 933 LOG.info(MessageFormat.format("{0} is set to false, not setting YARN restart properties", 934 HADOOP_YARN_KILL_CHILD_JOBS_ON_AMRESTART)); 935 } 936 937 String actionShareLibProperty = actionConf.get(ACTION_SHARELIB_FOR + getType()); 938 if (actionShareLibProperty != null) { 939 launcherJobConf.set(ACTION_SHARELIB_FOR + getType(), actionShareLibProperty); 940 } 941 setLibFilesArchives(context, actionXml, appPathRoot, launcherJobConf); 942 943 String jobName = launcherJobConf.get(HADOOP_JOB_NAME); 944 if (jobName == null || jobName.isEmpty()) { 945 jobName = XLog.format( 946 "oozie:launcher:T={0}:W={1}:A={2}:ID={3}", getType(), 947 context.getWorkflow().getAppName(), action.getName(), 948 context.getWorkflow().getId()); 949 launcherJobConf.setJobName(jobName); 950 } 951 952 // Inject Oozie job information if enabled. 953 injectJobInfo(launcherJobConf, actionConf, context, action); 954 955 injectLauncherCallback(context, launcherJobConf); 956 957 String jobId = context.getWorkflow().getId(); 958 String actionId = action.getId(); 959 Path actionDir = context.getActionDir(); 960 String recoveryId = context.getRecoveryId(); 961 962 // Getting the prepare XML from the action XML 963 Namespace ns = actionXml.getNamespace(); 964 Element prepareElement = actionXml.getChild("prepare", ns); 965 String prepareXML = ""; 966 if (prepareElement != null) { 967 if (prepareElement.getChildren().size() > 0) { 968 prepareXML = XmlUtils.prettyPrint(prepareElement).toString().trim(); 969 } 970 } 971 LauncherMapperHelper.setupLauncherInfo(launcherJobConf, jobId, actionId, actionDir, recoveryId, actionConf, 972 prepareXML); 973 974 // Set the launcher Main Class 975 LauncherMapperHelper.setupMainClass(launcherJobConf, getLauncherMain(launcherJobConf, actionXml)); 976 LauncherMapperHelper.setupLauncherURIHandlerConf(launcherJobConf); 977 978 LauncherMapperHelper.setupMaxOutputData(launcherJobConf, getMaxOutputData(actionConf)); 979 LauncherMapperHelper.setupMaxExternalStatsSize(launcherJobConf, maxExternalStatsSize); 980 LauncherMapperHelper.setupMaxFSGlob(launcherJobConf, maxFSGlobMax); 981 982 List<Element> list = actionXml.getChildren("arg", ns); 983 String[] args = new String[list.size()]; 984 for (int i = 0; i < list.size(); i++) { 985 args[i] = list.get(i).getTextTrim(); 986 } 987 LauncherMapperHelper.setupMainArguments(launcherJobConf, args); 988 // backward compatibility flag - see OOZIE-2872 989 if (ConfigurationService.getBoolean(LauncherMapper.CONF_OOZIE_NULL_ARGS_ALLOWED)) { 990 launcherJobConf.setBoolean(LauncherMapper.CONF_OOZIE_NULL_ARGS_ALLOWED, true); 991 } else { 992 launcherJobConf.setBoolean(LauncherMapper.CONF_OOZIE_NULL_ARGS_ALLOWED, false); 993 } 994 995 // Make mapred.child.java.opts and mapreduce.map.java.opts equal, but give values from the latter priority; also append 996 // <java-opt> and <java-opts> and give those highest priority 997 StringBuilder opts = new StringBuilder(launcherJobConf.get(HADOOP_CHILD_JAVA_OPTS, "")); 998 if (launcherJobConf.get(HADOOP_MAP_JAVA_OPTS) != null) { 999 opts.append(" ").append(launcherJobConf.get(HADOOP_MAP_JAVA_OPTS)); 1000 } 1001 List<Element> javaopts = actionXml.getChildren("java-opt", ns); 1002 for (Element opt: javaopts) { 1003 opts.append(" ").append(opt.getTextTrim()); 1004 } 1005 Element opt = actionXml.getChild("java-opts", ns); 1006 if (opt != null) { 1007 opts.append(" ").append(opt.getTextTrim()); 1008 } 1009 launcherJobConf.set(HADOOP_CHILD_JAVA_OPTS, opts.toString().trim()); 1010 launcherJobConf.set(HADOOP_MAP_JAVA_OPTS, opts.toString().trim()); 1011 1012 // setting for uber mode 1013 if (launcherJobConf.getBoolean(HADOOP_YARN_UBER_MODE, false)) { 1014 if (checkPropertiesToDisableUber(launcherJobConf)) { 1015 launcherJobConf.setBoolean(HADOOP_YARN_UBER_MODE, false); 1016 } 1017 else { 1018 updateConfForUberMode(launcherJobConf); 1019 } 1020 } 1021 updateConfForJavaTmpDir(launcherJobConf); 1022 1023 // properties from action that are needed by the launcher (e.g. QUEUE NAME, ACLs) 1024 // maybe we should add queue to the WF schema, below job-tracker 1025 actionConfToLauncherConf(actionConf, launcherJobConf); 1026 1027 checkAndSetupLauncherInputFormat(launcherJobConf); 1028 1029 return launcherJobConf; 1030 } 1031 catch (Exception ex) { 1032 throw convertException(ex); 1033 } 1034 } 1035 1036 private void checkAndSetupLauncherInputFormat(final JobConf launcherJobConf) throws ActionExecutorException { 1037 final Class inputFormatClass = launcherInputFormatClassLocator.locateOrGet(); 1038 LOG.debug("Launcher input format class is [{0}]", inputFormatClass.getName()); 1039 1040 if (!inputFormatClass.equals(OozieLauncherInputFormat.class) && FileInputFormat.class.isAssignableFrom(inputFormatClass)) { 1041 final String lastSharelibJarFullPath = getLastHDFSClasspathURI(launcherJobConf); 1042 LOG.debug("Adding placeholder input path [{0}] for launcher input format [{1}] " + 1043 "to work correctly as a FileInputFormat", 1044 lastSharelibJarFullPath, inputFormatClass); 1045 FileInputFormat.addInputPath(launcherJobConf, new Path(lastSharelibJarFullPath)); 1046 } 1047 } 1048 1049 @VisibleForTesting 1050 String getLastHDFSClasspathURI(final Configuration jobConf) throws ActionExecutorException { 1051 LOG.debug("Get last HDFS classpath URI"); 1052 1053 final String mrJobClasspathFiles = jobConf.get("mapreduce.job.classpath.files"); 1054 final String fsRootURI = jobConf.get(CommonConfigurationKeysPublic.FS_DEFAULT_NAME_KEY); 1055 1056 if (Strings.isEmpty(mrJobClasspathFiles) || Strings.isEmpty(fsRootURI)) { 1057 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "ALL", 1058 "Error submitting launcher. JobConf not configured correctly, cannot find last HDFS classpath URI."); 1059 } 1060 1061 String lastHDFSClasspathURI = null; 1062 1063 final String[] classpathURIs = mrJobClasspathFiles.split(","); 1064 int reverseIndex = classpathURIs.length - 1; 1065 1066 while (lastHDFSClasspathURI == null && reverseIndex >= 0) { 1067 final String classpathURI = classpathURIs[reverseIndex]; 1068 1069 if (classpathURI.contains(fsRootURI)) { 1070 lastHDFSClasspathURI = classpathURI.replace(fsRootURI, ""); 1071 } 1072 1073 reverseIndex--; 1074 } 1075 1076 if (Strings.isEmpty(lastHDFSClasspathURI)) { 1077 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "ALL", 1078 "Error submitting launcher. JobConf not configured correctly, there are no classpath entries stored on HDFS."); 1079 } 1080 1081 LOG.debug("Last HDFS classpath URI is [{0}]", lastHDFSClasspathURI); 1082 1083 return lastHDFSClasspathURI; 1084 } 1085 1086 @VisibleForTesting 1087 protected static int getMaxOutputData(Configuration actionConf) { 1088 String userMaxActionOutputLen = actionConf.get("oozie.action.max.output.data"); 1089 if (userMaxActionOutputLen != null) { 1090 Integer i = Ints.tryParse(userMaxActionOutputLen); 1091 return i != null ? i : maxActionOutputLen; 1092 } 1093 return maxActionOutputLen; 1094 } 1095 1096 private boolean checkPropertiesToDisableUber(Configuration launcherConf) { 1097 boolean disable = false; 1098 if (launcherConf.getBoolean(HADOOP_JOB_CLASSLOADER, false)) { 1099 disable = true; 1100 } 1101 else if (launcherConf.getBoolean(HADOOP_USER_CLASSPATH_FIRST, false)) { 1102 disable = true; 1103 } 1104 return disable; 1105 } 1106 1107 private void injectCallback(Context context, Configuration conf) { 1108 String callback = context.getCallbackUrl("$jobStatus"); 1109 if (conf.get("job.end.notification.url") != null) { 1110 LOG.warn("Overriding the action job end notification URI"); 1111 } 1112 conf.set("job.end.notification.url", callback); 1113 } 1114 1115 void injectActionCallback(Context context, Configuration actionConf) { 1116 injectCallback(context, actionConf); 1117 } 1118 1119 void injectLauncherCallback(Context context, Configuration launcherConf) { 1120 injectCallback(context, launcherConf); 1121 } 1122 1123 private void actionConfToLauncherConf(Configuration actionConf, JobConf launcherConf) { 1124 for (String name : SPECIAL_PROPERTIES) { 1125 if (actionConf.get(name) != null && launcherConf.get("oozie.launcher." + name) == null) { 1126 launcherConf.set(name, actionConf.get(name)); 1127 } 1128 } 1129 } 1130 1131 public void submitLauncher(FileSystem actionFs, Context context, WorkflowAction action) throws ActionExecutorException { 1132 JobClient jobClient = null; 1133 boolean exception = false; 1134 try { 1135 Path appPathRoot = new Path(context.getWorkflow().getAppPath()); 1136 1137 // app path could be a file 1138 if (actionFs.isFile(appPathRoot)) { 1139 appPathRoot = appPathRoot.getParent(); 1140 } 1141 1142 Element actionXml = XmlUtils.parseXml(action.getConf()); 1143 1144 // action job configuration 1145 Configuration actionConf = createBaseHadoopConf(context, actionXml); 1146 setupActionConf(actionConf, context, actionXml, appPathRoot); 1147 LOG.debug("Setting LibFilesArchives "); 1148 setLibFilesArchives(context, actionXml, appPathRoot, actionConf); 1149 1150 String jobName = actionConf.get(HADOOP_JOB_NAME); 1151 if (jobName == null || jobName.isEmpty()) { 1152 jobName = XLog.format("oozie:action:T={0}:W={1}:A={2}:ID={3}", 1153 getType(), context.getWorkflow().getAppName(), 1154 action.getName(), context.getWorkflow().getId()); 1155 actionConf.set(HADOOP_JOB_NAME, jobName); 1156 } 1157 1158 injectActionCallback(context, actionConf); 1159 1160 if(actionConf.get(ACL_MODIFY_JOB) == null || actionConf.get(ACL_MODIFY_JOB).trim().equals("")) { 1161 // ONLY in the case where user has not given the 1162 // modify-job ACL specifically 1163 if (context.getWorkflow().getAcl() != null) { 1164 // setting the group owning the Oozie job to allow anybody in that 1165 // group to modify the jobs. 1166 actionConf.set(ACL_MODIFY_JOB, context.getWorkflow().getAcl()); 1167 } 1168 } 1169 1170 // Setting the credential properties in launcher conf 1171 JobConf credentialsConf = null; 1172 HashMap<String, CredentialsProperties> credentialsProperties = setCredentialPropertyToActionConf(context, 1173 action, actionConf); 1174 if (credentialsProperties != null) { 1175 1176 // Adding if action need to set more credential tokens 1177 credentialsConf = new JobConf(false); 1178 XConfiguration.copy(actionConf, credentialsConf); 1179 setCredentialTokens(credentialsConf, context, action, credentialsProperties); 1180 1181 // insert conf to action conf from credentialsConf 1182 for (Entry<String, String> entry : credentialsConf) { 1183 if (actionConf.get(entry.getKey()) == null) { 1184 actionConf.set(entry.getKey(), entry.getValue()); 1185 } 1186 } 1187 } 1188 1189 JobConf launcherJobConf = createLauncherConf(actionFs, context, action, actionXml, actionConf); 1190 1191 LOG.debug("Creating Job Client for action " + action.getId()); 1192 jobClient = createJobClient(context, launcherJobConf); 1193 String launcherId = LauncherMapperHelper.getRecoveryId(launcherJobConf, context.getActionDir(), context 1194 .getRecoveryId()); 1195 boolean alreadyRunning = launcherId != null; 1196 RunningJob runningJob; 1197 1198 // if user-retry is on, always submit new launcher 1199 boolean isUserRetry = ((WorkflowActionBean)action).isUserRetry(); 1200 1201 if (alreadyRunning && !isUserRetry) { 1202 runningJob = jobClient.getJob(JobID.forName(launcherId)); 1203 if (runningJob == null) { 1204 String jobTracker = launcherJobConf.get(HADOOP_JOB_TRACKER); 1205 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "JA017", 1206 "unknown job [{0}@{1}], cannot recover", launcherId, jobTracker); 1207 } 1208 } 1209 else { 1210 LOG.debug("Submitting the job through Job Client for action " + action.getId()); 1211 1212 // setting up propagation of the delegation token. 1213 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 1214 Token<DelegationTokenIdentifier> mrdt = jobClient.getDelegationToken(has 1215 .getMRDelegationTokenRenewer(launcherJobConf)); 1216 launcherJobConf.getCredentials().addToken(HadoopAccessorService.MR_TOKEN_ALIAS, mrdt); 1217 1218 // insert credentials tokens to launcher job conf if needed 1219 if (needInjectCredentials() && credentialsConf != null) { 1220 for (Token<? extends TokenIdentifier> tk : credentialsConf.getCredentials().getAllTokens()) { 1221 Text fauxAlias = new Text(tk.getKind() + "_" + tk.getService()); 1222 LOG.debug("ADDING TOKEN: " + fauxAlias); 1223 launcherJobConf.getCredentials().addToken(fauxAlias, tk); 1224 } 1225 } 1226 else { 1227 LOG.info("No need to inject credentials."); 1228 } 1229 runningJob = jobClient.submitJob(launcherJobConf); 1230 if (runningJob == null) { 1231 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "JA017", 1232 "Error submitting launcher for action [{0}]", action.getId()); 1233 } 1234 launcherId = runningJob.getID().toString(); 1235 LOG.debug("After submission get the launcherId " + launcherId); 1236 } 1237 1238 String jobTracker = launcherJobConf.get(HADOOP_JOB_TRACKER); 1239 String consoleUrl = runningJob.getTrackingURL(); 1240 context.setStartData(launcherId, jobTracker, consoleUrl); 1241 } 1242 catch (Exception ex) { 1243 exception = true; 1244 throw convertException(ex); 1245 } 1246 finally { 1247 if (jobClient != null) { 1248 try { 1249 jobClient.close(); 1250 } 1251 catch (Exception e) { 1252 if (exception) { 1253 LOG.error("JobClient error: ", e); 1254 } 1255 else { 1256 throw convertException(e); 1257 } 1258 } 1259 } 1260 } 1261 } 1262 private boolean needInjectCredentials() { 1263 boolean methodExists = true; 1264 1265 Class klass; 1266 try { 1267 klass = Class.forName("org.apache.hadoop.mapred.JobConf"); 1268 klass.getMethod("getCredentials"); 1269 } 1270 catch (ClassNotFoundException ex) { 1271 methodExists = false; 1272 } 1273 catch (NoSuchMethodException ex) { 1274 methodExists = false; 1275 } 1276 1277 return methodExists; 1278 } 1279 1280 protected HashMap<String, CredentialsProperties> setCredentialPropertyToActionConf(Context context, 1281 WorkflowAction action, Configuration actionConf) throws Exception { 1282 HashMap<String, CredentialsProperties> credPropertiesMap = null; 1283 if (context != null && action != null) { 1284 if (!"true".equals(actionConf.get(OOZIE_CREDENTIALS_SKIP))) { 1285 XConfiguration wfJobConf = getWorkflowConf(context); 1286 if ("false".equals(actionConf.get(OOZIE_CREDENTIALS_SKIP)) || 1287 !wfJobConf.getBoolean(OOZIE_CREDENTIALS_SKIP, ConfigurationService.getBoolean(OOZIE_CREDENTIALS_SKIP))) { 1288 credPropertiesMap = getActionCredentialsProperties(context, action); 1289 if (credPropertiesMap != null) { 1290 for (String key : credPropertiesMap.keySet()) { 1291 CredentialsProperties prop = credPropertiesMap.get(key); 1292 if (prop != null) { 1293 LOG.debug("Credential Properties set for action : " + action.getId()); 1294 for (String property : prop.getProperties().keySet()) { 1295 actionConf.set(property, prop.getProperties().get(property)); 1296 LOG.debug("property : '" + property + "', value : '" + prop.getProperties().get(property) 1297 + "'"); 1298 } 1299 } 1300 } 1301 } else { 1302 LOG.warn("No credential properties found for action : " + action.getId() + ", cred : " + action.getCred()); 1303 } 1304 } else { 1305 LOG.info("Skipping credentials (" + OOZIE_CREDENTIALS_SKIP + "=true)"); 1306 } 1307 } else { 1308 LOG.info("Skipping credentials (" + OOZIE_CREDENTIALS_SKIP + "=true)"); 1309 } 1310 } else { 1311 LOG.warn("context or action is null"); 1312 } 1313 return credPropertiesMap; 1314 } 1315 1316 protected void setCredentialTokens(JobConf jobconf, Context context, WorkflowAction action, 1317 HashMap<String, CredentialsProperties> credPropertiesMap) throws Exception { 1318 1319 if (context != null && action != null && credPropertiesMap != null) { 1320 // Make sure we're logged into Kerberos; if not, or near expiration, it will relogin 1321 CredentialsProvider.ensureKerberosLogin(); 1322 for (Entry<String, CredentialsProperties> entry : credPropertiesMap.entrySet()) { 1323 String credName = entry.getKey(); 1324 CredentialsProperties credProps = entry.getValue(); 1325 if (credProps != null) { 1326 CredentialsProvider credProvider = new CredentialsProvider(credProps.getType()); 1327 Credentials credentialObject = credProvider.createCredentialObject(); 1328 if (credentialObject != null) { 1329 credentialObject.addtoJobConf(jobconf, credProps, context); 1330 LOG.debug("Retrieved Credential '" + credName + "' for action " + action.getId()); 1331 } 1332 else { 1333 LOG.debug("Credentials object is null for name= " + credName + ", type=" + credProps.getType()); 1334 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "JA020", 1335 "Could not load credentials of type [{0}] with name [{1}]]; perhaps it was not defined" 1336 + " in oozie-site.xml?", credProps.getType(), credName); 1337 } 1338 } 1339 } 1340 } 1341 1342 } 1343 1344 protected HashMap<String, CredentialsProperties> getActionCredentialsProperties(Context context, 1345 WorkflowAction action) throws Exception { 1346 HashMap<String, CredentialsProperties> props = new HashMap<String, CredentialsProperties>(); 1347 if (context != null && action != null) { 1348 String credsInAction = action.getCred(); 1349 if (credsInAction != null) { 1350 LOG.debug("Get credential '" + credsInAction + "' properties for action : " + action.getId()); 1351 String[] credNames = credsInAction.split(","); 1352 for (String credName : credNames) { 1353 CredentialsProperties credProps = getCredProperties(context, credName); 1354 props.put(credName, credProps); 1355 } 1356 } 1357 } 1358 else { 1359 LOG.warn("context or action is null"); 1360 } 1361 return props; 1362 } 1363 1364 @SuppressWarnings("unchecked") 1365 protected CredentialsProperties getCredProperties(Context context, String credName) 1366 throws Exception { 1367 CredentialsProperties credProp = null; 1368 String workflowXml = ((WorkflowJobBean) context.getWorkflow()).getWorkflowInstance().getApp().getDefinition(); 1369 XConfiguration wfjobConf = getWorkflowConf(context); 1370 Element elementJob = XmlUtils.parseXml(workflowXml); 1371 Element credentials = elementJob.getChild("credentials", elementJob.getNamespace()); 1372 if (credentials != null) { 1373 for (Element credential : (List<Element>) credentials.getChildren("credential", credentials.getNamespace())) { 1374 String name = credential.getAttributeValue("name"); 1375 String type = credential.getAttributeValue("type"); 1376 LOG.debug("getCredProperties: Name: " + name + ", Type: " + type); 1377 if (name.equalsIgnoreCase(credName)) { 1378 credProp = new CredentialsProperties(name, type); 1379 for (Element property : (List<Element>) credential.getChildren("property", 1380 credential.getNamespace())) { 1381 String propertyName = property.getChildText("name", property.getNamespace()); 1382 String propertyValue = property.getChildText("value", property.getNamespace()); 1383 ELEvaluator eval = new ELEvaluator(); 1384 for (Map.Entry<String, String> entry : wfjobConf) { 1385 eval.setVariable(entry.getKey(), entry.getValue().trim()); 1386 } 1387 propertyName = eval.evaluate(propertyName, String.class); 1388 propertyValue = eval.evaluate(propertyValue, String.class); 1389 1390 credProp.getProperties().put(propertyName, propertyValue); 1391 LOG.debug("getCredProperties: Properties name :'" + propertyName + "', Value : '" 1392 + propertyValue + "'"); 1393 } 1394 } 1395 } 1396 } else { 1397 LOG.debug("credentials is null for the action"); 1398 } 1399 return credProp; 1400 } 1401 1402 @Override 1403 public void start(Context context, WorkflowAction action) throws ActionExecutorException { 1404 LogUtils.setLogInfo(action); 1405 try { 1406 LOG.debug("Starting action " + action.getId() + " getting Action File System"); 1407 FileSystem actionFs = context.getAppFileSystem(); 1408 LOG.debug("Preparing action Dir through copying " + context.getActionDir()); 1409 prepareActionDir(actionFs, context); 1410 LOG.debug("Action Dir is ready. Submitting the action "); 1411 submitLauncher(actionFs, context, action); 1412 LOG.debug("Action submit completed. Performing check "); 1413 check(context, action); 1414 LOG.debug("Action check is done after submission"); 1415 } 1416 catch (Exception ex) { 1417 throw convertException(ex); 1418 } 1419 } 1420 1421 @Override 1422 public void end(Context context, WorkflowAction action) throws ActionExecutorException { 1423 try { 1424 String externalStatus = action.getExternalStatus(); 1425 WorkflowAction.Status status = externalStatus.equals(SUCCEEDED) ? WorkflowAction.Status.OK 1426 : WorkflowAction.Status.ERROR; 1427 context.setEndData(status, getActionSignal(status)); 1428 } 1429 catch (Exception ex) { 1430 throw convertException(ex); 1431 } 1432 finally { 1433 try { 1434 FileSystem actionFs = context.getAppFileSystem(); 1435 cleanUpActionDir(actionFs, context); 1436 } 1437 catch (Exception ex) { 1438 throw convertException(ex); 1439 } 1440 } 1441 } 1442 1443 /** 1444 * Create job client object 1445 * 1446 * @param context 1447 * @param jobConf 1448 * @return JobClient 1449 * @throws HadoopAccessorException 1450 */ 1451 protected JobClient createJobClient(Context context, JobConf jobConf) throws HadoopAccessorException { 1452 String user = context.getWorkflow().getUser(); 1453 String group = context.getWorkflow().getGroup(); 1454 return Services.get().get(HadoopAccessorService.class).createJobClient(user, jobConf); 1455 } 1456 1457 protected RunningJob getRunningJob(Context context, WorkflowAction action, JobClient jobClient) throws Exception{ 1458 RunningJob runningJob = jobClient.getJob(JobID.forName(action.getExternalId())); 1459 return runningJob; 1460 } 1461 1462 /** 1463 * Useful for overriding in actions that do subsequent job runs 1464 * such as the MapReduce Action, where the launcher job is not the 1465 * actual job that then gets monitored. 1466 */ 1467 protected String getActualExternalId(WorkflowAction action) { 1468 return action.getExternalId(); 1469 } 1470 1471 @Override 1472 public void check(Context context, WorkflowAction action) throws ActionExecutorException { 1473 JobClient jobClient = null; 1474 boolean exception = false; 1475 LogUtils.setLogInfo(action); 1476 try { 1477 Element actionXml = XmlUtils.parseXml(action.getConf()); 1478 FileSystem actionFs = context.getAppFileSystem(); 1479 JobConf jobConf = createBaseHadoopConf(context, actionXml); 1480 jobClient = createJobClient(context, jobConf); 1481 RunningJob runningJob = getRunningJob(context, action, jobClient); 1482 if (runningJob == null) { 1483 context.setExecutionData(FAILED, null); 1484 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "JA017", 1485 "Could not lookup launched hadoop Job ID [{0}] which was associated with " + 1486 " action [{1}]. Failing this action!", getActualExternalId(action), action.getId()); 1487 } 1488 if (runningJob.isComplete()) { 1489 Path actionDir = context.getActionDir(); 1490 String newId = null; 1491 // load sequence file into object 1492 Map<String, String> actionData = LauncherMapperHelper.getActionData(actionFs, actionDir, jobConf); 1493 if (actionData.containsKey(LauncherMapper.ACTION_DATA_NEW_ID)) { 1494 newId = actionData.get(LauncherMapper.ACTION_DATA_NEW_ID); 1495 String launcherId = action.getExternalId(); 1496 runningJob = jobClient.getJob(JobID.forName(newId)); 1497 if (runningJob == null) { 1498 context.setExternalStatus(FAILED); 1499 throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "JA017", 1500 "Unknown hadoop job [{0}] associated with action [{1}]. Failing this action!", newId, 1501 action.getId()); 1502 } 1503 context.setExternalChildIDs(newId); 1504 LOG.info(XLog.STD, "External ID swap, old ID [{0}] new ID [{1}]", launcherId, 1505 newId); 1506 } 1507 else { 1508 String externalIDs = actionData.get(LauncherMapper.ACTION_DATA_EXTERNAL_CHILD_IDS); 1509 if (externalIDs != null) { 1510 context.setExternalChildIDs(externalIDs); 1511 LOG.info(XLog.STD, "Hadoop Jobs launched : [{0}]", externalIDs); 1512 } 1513 else if (LauncherMapperHelper.hasOutputData(actionData)) { 1514 // Load stored Hadoop jobs ids and promote them as external child ids 1515 // This is for jobs launched with older release during upgrade to Oozie 4.3 1516 Properties props = PropertiesUtils.stringToProperties(actionData 1517 .get(LauncherMapper.ACTION_DATA_OUTPUT_PROPS)); 1518 if (props.get(LauncherMain.HADOOP_JOBS) != null) { 1519 externalIDs = (String) props.get(LauncherMain.HADOOP_JOBS); 1520 context.setExternalChildIDs(externalIDs); 1521 LOG.info(XLog.STD, "Hadoop Jobs launched : [{0}]", externalIDs); 1522 } 1523 } 1524 } 1525 if (runningJob.isComplete()) { 1526 // fetching action output and stats for the Map-Reduce action. 1527 if (newId != null) { 1528 actionData = LauncherMapperHelper.getActionData(actionFs, context.getActionDir(), jobConf); 1529 } 1530 LOG.info(XLog.STD, "action completed, external ID [{0}]", 1531 action.getExternalId()); 1532 if (LauncherMapperHelper.isMainSuccessful(runningJob)) { 1533 if (getCaptureOutput(action) && LauncherMapperHelper.hasOutputData(actionData)) { 1534 context.setExecutionData(SUCCEEDED, PropertiesUtils.stringToProperties(actionData 1535 .get(LauncherMapper.ACTION_DATA_OUTPUT_PROPS))); 1536 LOG.info(XLog.STD, "action produced output"); 1537 } 1538 else { 1539 context.setExecutionData(SUCCEEDED, null); 1540 } 1541 if (LauncherMapperHelper.hasStatsData(actionData)) { 1542 context.setExecutionStats(actionData.get(LauncherMapper.ACTION_DATA_STATS)); 1543 LOG.info(XLog.STD, "action produced stats"); 1544 } 1545 getActionData(actionFs, runningJob, action, context); 1546 } 1547 else { 1548 String errorReason; 1549 if (actionData.containsKey(LauncherMapper.ACTION_DATA_ERROR_PROPS)) { 1550 Properties props = PropertiesUtils.stringToProperties(actionData 1551 .get(LauncherMapper.ACTION_DATA_ERROR_PROPS)); 1552 String errorCode = props.getProperty("error.code"); 1553 if ("0".equals(errorCode)) { 1554 errorCode = "JA018"; 1555 } 1556 if ("-1".equals(errorCode)) { 1557 errorCode = "JA019"; 1558 } 1559 errorReason = props.getProperty("error.reason"); 1560 LOG.warn("Launcher ERROR, reason: {0}", errorReason); 1561 String exMsg = props.getProperty("exception.message"); 1562 String errorInfo = (exMsg != null) ? exMsg : errorReason; 1563 context.setErrorInfo(errorCode, errorInfo); 1564 String exStackTrace = props.getProperty("exception.stacktrace"); 1565 if (exMsg != null) { 1566 LOG.warn("Launcher exception: {0}{E}{1}", exMsg, exStackTrace); 1567 } 1568 } 1569 else { 1570 errorReason = XLog.format("LauncherMapper died, check Hadoop LOG for job [{0}:{1}]", action 1571 .getTrackerUri(), action.getExternalId()); 1572 LOG.warn(errorReason); 1573 } 1574 context.setExecutionData(FAILED_KILLED, null); 1575 } 1576 } 1577 else { 1578 context.setExternalStatus("RUNNING"); 1579 LOG.info(XLog.STD, "checking action, hadoop job ID [{0}] status [RUNNING]", 1580 runningJob.getID()); 1581 } 1582 } 1583 else { 1584 context.setExternalStatus("RUNNING"); 1585 LOG.info(XLog.STD, "checking action, hadoop job ID [{0}] status [RUNNING]", 1586 runningJob.getID()); 1587 } 1588 } 1589 catch (Exception ex) { 1590 LOG.warn("Exception in check(). Message[{0}]", ex.getMessage(), ex); 1591 exception = true; 1592 throw convertException(ex); 1593 } 1594 finally { 1595 if (jobClient != null) { 1596 try { 1597 jobClient.close(); 1598 } 1599 catch (Exception e) { 1600 if (exception) { 1601 LOG.error("JobClient error: ", e); 1602 } 1603 else { 1604 throw convertException(e); 1605 } 1606 } 1607 } 1608 } 1609 } 1610 1611 /** 1612 * Get the output data of an action. Subclasses should override this method 1613 * to get action specific output data. 1614 * 1615 * @param actionFs the FileSystem object 1616 * @param runningJob the runningJob 1617 * @param action the Workflow action 1618 * @param context executor context 1619 * 1620 */ 1621 protected void getActionData(FileSystem actionFs, RunningJob runningJob, WorkflowAction action, Context context) 1622 throws HadoopAccessorException, JDOMException, IOException, URISyntaxException { 1623 } 1624 1625 protected boolean getCaptureOutput(WorkflowAction action) throws JDOMException { 1626 Element eConf = XmlUtils.parseXml(action.getConf()); 1627 Namespace ns = eConf.getNamespace(); 1628 Element captureOutput = eConf.getChild("capture-output", ns); 1629 return captureOutput != null; 1630 } 1631 1632 @Override 1633 public void kill(Context context, WorkflowAction action) throws ActionExecutorException { 1634 JobClient jobClient = null; 1635 boolean exception = false; 1636 try { 1637 Element actionXml = XmlUtils.parseXml(action.getConf()); 1638 final JobConf jobConf = createBaseHadoopConf(context, actionXml); 1639 WorkflowJob wfJob = context.getWorkflow(); 1640 Configuration conf = null; 1641 if ( wfJob.getConf() != null ) { 1642 conf = new XConfiguration(new StringReader(wfJob.getConf())); 1643 } 1644 String launcherTag = LauncherMapperHelper.getActionYarnTag(conf, wfJob.getParentId(), action); 1645 jobConf.set(LauncherMainHadoopUtils.CHILD_MAPREDUCE_JOB_TAGS, LauncherMapperHelper.getTag(launcherTag)); 1646 jobConf.set(LauncherMainHadoopUtils.OOZIE_JOB_LAUNCH_TIME, Long.toString(action.getStartTime().getTime())); 1647 UserGroupInformation ugi = Services.get().get(UserGroupInformationService.class) 1648 .getProxyUser(context.getWorkflow().getUser()); 1649 ugi.doAs(new PrivilegedExceptionAction<Void>() { 1650 @Override 1651 public Void run() throws Exception { 1652 LauncherMainHadoopUtils.killChildYarnJobs(jobConf); 1653 return null; 1654 } 1655 }); 1656 jobClient = createJobClient(context, jobConf); 1657 RunningJob runningJob = getRunningJob(context, action, jobClient); 1658 if (runningJob != null) { 1659 runningJob.killJob(); 1660 } 1661 context.setExternalStatus(KILLED); 1662 context.setExecutionData(KILLED, null); 1663 } 1664 catch (Exception ex) { 1665 exception = true; 1666 throw convertException(ex); 1667 } 1668 finally { 1669 try { 1670 FileSystem actionFs = context.getAppFileSystem(); 1671 cleanUpActionDir(actionFs, context); 1672 if (jobClient != null) { 1673 jobClient.close(); 1674 } 1675 } 1676 catch (Exception ex) { 1677 if (exception) { 1678 LOG.error("Error: ", ex); 1679 } 1680 else { 1681 throw convertException(ex); 1682 } 1683 } 1684 } 1685 } 1686 1687 private static Set<String> FINAL_STATUS = new HashSet<String>(); 1688 1689 static { 1690 FINAL_STATUS.add(SUCCEEDED); 1691 FINAL_STATUS.add(KILLED); 1692 FINAL_STATUS.add(FAILED); 1693 FINAL_STATUS.add(FAILED_KILLED); 1694 } 1695 1696 @Override 1697 public boolean isCompleted(String externalStatus) { 1698 return FINAL_STATUS.contains(externalStatus); 1699 } 1700 1701 1702 /** 1703 * Return the sharelib names for the action. 1704 * <p/> 1705 * If <code>NULL</code> or empty, it means that the action does not use the action 1706 * sharelib. 1707 * <p/> 1708 * If a non-empty string, i.e. <code>foo</code>, it means the action uses the 1709 * action sharelib sub-directory <code>foo</code> and all JARs in the sharelib 1710 * <code>foo</code> directory will be in the action classpath. Multiple sharelib 1711 * sub-directories can be specified as a comma separated list. 1712 * <p/> 1713 * The resolution is done using the following precedence order: 1714 * <ul> 1715 * <li><b>action.sharelib.for.#ACTIONTYPE#</b> in the action configuration</li> 1716 * <li><b>action.sharelib.for.#ACTIONTYPE#</b> in the job configuration</li> 1717 * <li><b>action.sharelib.for.#ACTIONTYPE#</b> in the oozie configuration</li> 1718 * <li>Action Executor <code>getDefaultShareLibName()</code> method</li> 1719 * </ul> 1720 * 1721 * 1722 * @param context executor context. 1723 * @param actionXml 1724 * @param conf action configuration. 1725 * @return the action sharelib names. 1726 */ 1727 protected String[] getShareLibNames(Context context, Element actionXml, Configuration conf) { 1728 String[] names = conf.getStrings(ACTION_SHARELIB_FOR + getType()); 1729 if (names == null || names.length == 0) { 1730 try { 1731 XConfiguration jobConf = getWorkflowConf(context); 1732 names = jobConf.getStrings(ACTION_SHARELIB_FOR + getType()); 1733 if (names == null || names.length == 0) { 1734 names = Services.get().getConf().getStrings(ACTION_SHARELIB_FOR + getType()); 1735 if (names == null || names.length == 0) { 1736 String name = getDefaultShareLibName(actionXml); 1737 if (name != null) { 1738 names = new String[] { name }; 1739 } 1740 } 1741 } 1742 } 1743 catch (IOException ex) { 1744 throw new RuntimeException("It cannot happen, " + ex.toString(), ex); 1745 } 1746 } 1747 return names; 1748 } 1749 1750 private final static String ACTION_SHARELIB_FOR = "oozie.action.sharelib.for."; 1751 1752 1753 /** 1754 * Returns the default sharelib name for the action if any. 1755 * 1756 * @param actionXml the action XML fragment. 1757 * @return the sharelib name for the action, <code>NULL</code> if none. 1758 */ 1759 protected String getDefaultShareLibName(Element actionXml) { 1760 return null; 1761 } 1762 1763 public String[] getShareLibFilesForActionConf() { 1764 return null; 1765 } 1766 1767 /** 1768 * Sets some data for the action on completion 1769 * 1770 * @param context executor context 1771 * @param actionFs the FileSystem object 1772 */ 1773 protected void setActionCompletionData(Context context, FileSystem actionFs) throws IOException, 1774 HadoopAccessorException, URISyntaxException { 1775 } 1776 1777 private void injectJobInfo(JobConf launcherJobConf, Configuration actionConf, Context context, WorkflowAction action) { 1778 if (OozieJobInfo.isJobInfoEnabled()) { 1779 try { 1780 OozieJobInfo jobInfo = new OozieJobInfo(actionConf, context, action); 1781 String jobInfoStr = jobInfo.getJobInfo(); 1782 launcherJobConf.set(OozieJobInfo.JOB_INFO_KEY, jobInfoStr + "launcher=true"); 1783 actionConf.set(OozieJobInfo.JOB_INFO_KEY, jobInfoStr + "launcher=false"); 1784 } 1785 catch (Exception e) { 1786 // Just job info, should not impact the execution. 1787 LOG.error("Error while populating job info", e); 1788 } 1789 } 1790 } 1791 1792 @Override 1793 public boolean requiresNameNodeJobTracker() { 1794 return true; 1795 } 1796 1797 @Override 1798 public boolean supportsConfigurationJobXML() { 1799 return true; 1800 } 1801 1802 private XConfiguration getWorkflowConf(Context context) throws IOException { 1803 if (workflowConf == null) { 1804 workflowConf = new XConfiguration(new StringReader(context.getWorkflow().getConf())); 1805 } 1806 return workflowConf; 1807 1808 } 1809 1810 private String getActionTypeLauncherPrefix() { 1811 return "oozie.action." + getType() + ".launcher."; 1812 } 1813}