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