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