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