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.command.coord; 020 021import java.io.IOException; 022import java.io.InputStreamReader; 023import java.io.Reader; 024import java.io.StringReader; 025import java.io.StringWriter; 026import java.net.URI; 027import java.net.URISyntaxException; 028import java.util.ArrayList; 029import java.util.Calendar; 030import java.util.Date; 031import java.util.HashMap; 032import java.util.HashSet; 033import java.util.Iterator; 034import java.util.List; 035import java.util.Map; 036import java.util.Set; 037import java.util.TreeSet; 038 039import javax.xml.transform.stream.StreamSource; 040import javax.xml.validation.Validator; 041 042import org.apache.hadoop.conf.Configuration; 043import org.apache.hadoop.fs.FileSystem; 044import org.apache.hadoop.fs.Path; 045import org.apache.oozie.CoordinatorJobBean; 046import org.apache.oozie.ErrorCode; 047import org.apache.oozie.client.CoordinatorJob; 048import org.apache.oozie.client.Job; 049import org.apache.oozie.client.OozieClient; 050import org.apache.oozie.client.CoordinatorJob.Execution; 051import org.apache.oozie.command.CommandException; 052import org.apache.oozie.command.SubmitTransitionXCommand; 053import org.apache.oozie.command.bundle.BundleStatusUpdateXCommand; 054import org.apache.oozie.coord.CoordELEvaluator; 055import org.apache.oozie.coord.CoordELFunctions; 056import org.apache.oozie.coord.CoordinatorJobException; 057import org.apache.oozie.coord.TimeUnit; 058import org.apache.oozie.executor.jpa.CoordJobQueryExecutor; 059import org.apache.oozie.executor.jpa.JPAExecutorException; 060import org.apache.oozie.service.CoordMaterializeTriggerService; 061import org.apache.oozie.service.ConfigurationService; 062import org.apache.oozie.service.HadoopAccessorException; 063import org.apache.oozie.service.HadoopAccessorService; 064import org.apache.oozie.service.JPAService; 065import org.apache.oozie.service.SchemaService; 066import org.apache.oozie.service.Service; 067import org.apache.oozie.service.Services; 068import org.apache.oozie.service.UUIDService; 069import org.apache.oozie.service.SchemaService.SchemaName; 070import org.apache.oozie.service.UUIDService.ApplicationType; 071import org.apache.oozie.util.ConfigUtils; 072import org.apache.oozie.util.DateUtils; 073import org.apache.oozie.util.ELEvaluator; 074import org.apache.oozie.util.ELUtils; 075import org.apache.oozie.util.IOUtils; 076import org.apache.oozie.util.InstrumentUtils; 077import org.apache.oozie.util.LogUtils; 078import org.apache.oozie.util.ParamChecker; 079import org.apache.oozie.util.ParameterVerifier; 080import org.apache.oozie.util.ParameterVerifierException; 081import org.apache.oozie.util.PropertiesUtils; 082import org.apache.oozie.util.XConfiguration; 083import org.apache.oozie.util.XmlUtils; 084import org.jdom.Attribute; 085import org.jdom.Element; 086import org.jdom.JDOMException; 087import org.jdom.Namespace; 088import org.xml.sax.SAXException; 089 090/** 091 * This class provides the functionalities to resolve a coordinator job XML and write the job information into a DB 092 * table. 093 * <p/> 094 * Specifically it performs the following functions: 1. Resolve all the variables or properties using job 095 * configurations. 2. Insert all datasets definition as part of the <data-in> and <data-out> tags. 3. Validate the XML 096 * at runtime. 097 */ 098public class CoordSubmitXCommand extends SubmitTransitionXCommand { 099 100 protected Configuration conf; 101 private final String bundleId; 102 private final String coordName; 103 protected boolean dryrun; 104 protected JPAService jpaService = null; 105 private CoordinatorJob.Status prevStatus = CoordinatorJob.Status.PREP; 106 107 public static final String CONFIG_DEFAULT = "coord-config-default.xml"; 108 public static final String COORDINATOR_XML_FILE = "coordinator.xml"; 109 public final String COORD_INPUT_EVENTS ="input-events"; 110 public final String COORD_OUTPUT_EVENTS = "output-events"; 111 public final String COORD_INPUT_EVENTS_DATA_IN ="data-in"; 112 public final String COORD_OUTPUT_EVENTS_DATA_OUT = "data-out"; 113 114 private static final Set<String> DISALLOWED_USER_PROPERTIES = new HashSet<String>(); 115 private static final Set<String> DISALLOWED_DEFAULT_PROPERTIES = new HashSet<String>(); 116 117 protected CoordinatorJobBean coordJob = null; 118 /** 119 * Default timeout for normal jobs, in minutes, after which coordinator input check will timeout 120 */ 121 public static final String CONF_DEFAULT_TIMEOUT_NORMAL = Service.CONF_PREFIX + "coord.normal.default.timeout"; 122 123 public static final String CONF_DEFAULT_CONCURRENCY = Service.CONF_PREFIX + "coord.default.concurrency"; 124 125 public static final String CONF_DEFAULT_THROTTLE = Service.CONF_PREFIX + "coord.default.throttle"; 126 127 public static final String CONF_MAT_THROTTLING_FACTOR = Service.CONF_PREFIX 128 + "coord.materialization.throttling.factor"; 129 130 /** 131 * Default MAX timeout in minutes, after which coordinator input check will timeout 132 */ 133 public static final String CONF_DEFAULT_MAX_TIMEOUT = Service.CONF_PREFIX + "coord.default.max.timeout"; 134 135 public static final String CONF_QUEUE_SIZE = Service.CONF_PREFIX + "CallableQueueService.queue.size"; 136 137 public static final String CONF_CHECK_MAX_FREQUENCY = Service.CONF_PREFIX + "coord.check.maximum.frequency"; 138 139 private ELEvaluator evalFreq = null; 140 private ELEvaluator evalNofuncs = null; 141 private ELEvaluator evalData = null; 142 private ELEvaluator evalInst = null; 143 private ELEvaluator evalAction = null; 144 private ELEvaluator evalSla = null; 145 private ELEvaluator evalTimeout = null; 146 private ELEvaluator evalInitialInstance = null; 147 148 static { 149 String[] badUserProps = { PropertiesUtils.YEAR, PropertiesUtils.MONTH, PropertiesUtils.DAY, 150 PropertiesUtils.HOUR, PropertiesUtils.MINUTE, PropertiesUtils.DAYS, PropertiesUtils.HOURS, 151 PropertiesUtils.MINUTES, PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, 152 PropertiesUtils.TB, PropertiesUtils.PB, PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, 153 PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN, PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS }; 154 PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_USER_PROPERTIES); 155 } 156 157 /** 158 * Constructor to create the Coordinator Submit Command. 159 * 160 * @param conf : Configuration for Coordinator job 161 */ 162 public CoordSubmitXCommand(Configuration conf) { 163 super("coord_submit", "coord_submit", 1); 164 this.conf = ParamChecker.notNull(conf, "conf"); 165 this.bundleId = null; 166 this.coordName = null; 167 } 168 169 /** 170 * Constructor to create the Coordinator Submit Command by bundle job. 171 * 172 * @param conf : Configuration for Coordinator job 173 * @param bundleId : bundle id 174 * @param coordName : coord name 175 */ 176 public CoordSubmitXCommand(Configuration conf, String bundleId, String coordName) { 177 super("coord_submit", "coord_submit", 1); 178 this.conf = ParamChecker.notNull(conf, "conf"); 179 this.bundleId = ParamChecker.notEmpty(bundleId, "bundleId"); 180 this.coordName = ParamChecker.notEmpty(coordName, "coordName"); 181 } 182 183 /** 184 * Constructor to create the Coordinator Submit Command. 185 * 186 * @param dryrun : if dryrun 187 * @param conf : Configuration for Coordinator job 188 */ 189 public CoordSubmitXCommand(boolean dryrun, Configuration conf) { 190 this(conf); 191 this.dryrun = dryrun; 192 } 193 194 /* (non-Javadoc) 195 * @see org.apache.oozie.command.XCommand#execute() 196 */ 197 @Override 198 protected String submit() throws CommandException { 199 LOG.info("STARTED Coordinator Submit"); 200 String jobId = submitJob(); 201 LOG.info("ENDED Coordinator Submit jobId=" + jobId); 202 return jobId; 203 } 204 205 protected String submitJob() throws CommandException { 206 String jobId = null; 207 InstrumentUtils.incrJobCounter(getName(), 1, getInstrumentation()); 208 209 boolean exceptionOccured = false; 210 try { 211 mergeDefaultConfig(); 212 213 String appXml = readAndValidateXml(); 214 coordJob.setOrigJobXml(appXml); 215 LOG.debug("jobXml after initial validation " + XmlUtils.prettyPrint(appXml).toString()); 216 217 Element eXml = XmlUtils.parseXml(appXml); 218 219 String appNamespace = readAppNamespace(eXml); 220 coordJob.setAppNamespace(appNamespace); 221 222 ParameterVerifier.verifyParameters(conf, eXml); 223 224 appXml = XmlUtils.removeComments(appXml); 225 initEvaluators(); 226 Element eJob = basicResolveAndIncludeDS(appXml, conf, coordJob); 227 228 validateCoordinatorJob(); 229 230 // checking if the coordinator application data input/output events 231 // specify multiple data instance values in erroneous manner 232 checkMultipleTimeInstances(eJob, COORD_INPUT_EVENTS, COORD_INPUT_EVENTS_DATA_IN); 233 checkMultipleTimeInstances(eJob, COORD_OUTPUT_EVENTS, COORD_OUTPUT_EVENTS_DATA_OUT); 234 235 LOG.debug("jobXml after all validation " + XmlUtils.prettyPrint(eJob).toString()); 236 237 jobId = storeToDB(appXml, eJob, coordJob); 238 // log job info for coordinator job 239 LogUtils.setLogInfo(coordJob); 240 241 if (!dryrun) { 242 queueMaterializeTransitionXCommand(jobId); 243 } 244 else { 245 return getDryRun(coordJob); 246 } 247 } 248 catch (JDOMException jex) { 249 exceptionOccured = true; 250 LOG.warn("ERROR: ", jex); 251 throw new CommandException(ErrorCode.E0700, jex.getMessage(), jex); 252 } 253 catch (CoordinatorJobException cex) { 254 exceptionOccured = true; 255 LOG.warn("ERROR: ", cex); 256 throw new CommandException(cex); 257 } 258 catch (ParameterVerifierException pex) { 259 exceptionOccured = true; 260 LOG.warn("ERROR: ", pex); 261 throw new CommandException(pex); 262 } 263 catch (IllegalArgumentException iex) { 264 exceptionOccured = true; 265 LOG.warn("ERROR: ", iex); 266 throw new CommandException(ErrorCode.E1003, iex.getMessage(), iex); 267 } 268 catch (Exception ex) { 269 exceptionOccured = true; 270 LOG.warn("ERROR: ", ex); 271 throw new CommandException(ErrorCode.E0803, ex.getMessage(), ex); 272 } 273 finally { 274 if (exceptionOccured) { 275 if (coordJob.getId() == null || coordJob.getId().equalsIgnoreCase("")) { 276 coordJob.setStatus(CoordinatorJob.Status.FAILED); 277 coordJob.resetPending(); 278 } 279 } 280 } 281 return jobId; 282 } 283 284 /** 285 * Gets the dryrun output. 286 * 287 * @param coordJob the coordinatorJobBean 288 * @return the dry run 289 * @throws Exception the exception 290 */ 291 protected String getDryRun(CoordinatorJobBean coordJob) throws Exception{ 292 int materializationWindow = ConfigurationService 293 .getInt(CoordMaterializeTriggerService.CONF_MATERIALIZATION_WINDOW); 294 Date startTime = coordJob.getStartTime(); 295 long startTimeMilli = startTime.getTime(); 296 long endTimeMilli = startTimeMilli + (materializationWindow * 1000); 297 Date jobEndTime = coordJob.getEndTime(); 298 Date endTime = new Date(endTimeMilli); 299 if (endTime.compareTo(jobEndTime) > 0) { 300 endTime = jobEndTime; 301 } 302 String jobId = coordJob.getId(); 303 LOG.info("[" + jobId + "]: Update status to RUNNING"); 304 coordJob.setStatus(Job.Status.RUNNING); 305 coordJob.setPending(); 306 Configuration jobConf = null; 307 try { 308 jobConf = new XConfiguration(new StringReader(coordJob.getConf())); 309 } 310 catch (IOException e1) { 311 LOG.warn("Configuration parse error. read from DB :" + coordJob.getConf(), e1); 312 } 313 String action = new CoordMaterializeTransitionXCommand(coordJob, materializationWindow, startTime, 314 endTime).materializeActions(true); 315 String output = coordJob.getJobXml() + System.getProperty("line.separator") 316 + "***actions for instance***" + action; 317 return output; 318 } 319 320 /** 321 * Queue MaterializeTransitionXCommand 322 */ 323 protected void queueMaterializeTransitionXCommand(String jobId) { 324 int materializationWindow = ConfigurationService 325 .getInt(CoordMaterializeTriggerService.CONF_MATERIALIZATION_WINDOW); 326 queue(new CoordMaterializeTransitionXCommand(jobId, materializationWindow), 100); 327 } 328 329 /** 330 * Method that validates values in the definition for correctness. Placeholder to add more. 331 */ 332 private void validateCoordinatorJob() throws Exception { 333 // check if startTime < endTime 334 if (!coordJob.getStartTime().before(coordJob.getEndTime())) { 335 throw new IllegalArgumentException("Coordinator Start Time must be earlier than End Time."); 336 } 337 338 try { 339 // Check if a coord job with cron frequency will materialize actions 340 int freq = Integer.parseInt(coordJob.getFrequency()); 341 342 // Check if the frequency is faster than 5 min if enabled 343 if (ConfigurationService.getBoolean(CONF_CHECK_MAX_FREQUENCY)) { 344 CoordinatorJob.Timeunit unit = coordJob.getTimeUnit(); 345 if (freq == 0 || (freq < 5 && unit == CoordinatorJob.Timeunit.MINUTE)) { 346 throw new IllegalArgumentException("Coordinator job with frequency [" + freq + 347 "] minutes is faster than allowed maximum of 5 minutes (" 348 + CONF_CHECK_MAX_FREQUENCY + " is set to true)"); 349 } 350 } 351 } catch (NumberFormatException e) { 352 Date start = coordJob.getStartTime(); 353 Calendar cal = Calendar.getInstance(); 354 cal.setTime(start); 355 cal.add(Calendar.MINUTE, -1); 356 start = cal.getTime(); 357 358 Date nextTime = CoordCommandUtils.getNextValidActionTimeForCronFrequency(start, coordJob); 359 if (nextTime == null) { 360 throw new IllegalArgumentException("Invalid coordinator cron frequency: " + coordJob.getFrequency()); 361 } 362 if (!nextTime.before(coordJob.getEndTime())) { 363 throw new IllegalArgumentException("Coordinator job with frequency '" + 364 coordJob.getFrequency() + "' materializes no actions between start and end time."); 365 } 366 } 367 } 368 369 /* 370 * Check against multiple data instance values inside a single <instance> <start-instance> or <end-instance> tag 371 * If found, the job is not submitted and user is informed to correct the error, instead of defaulting to the first instance value in the list 372 */ 373 private void checkMultipleTimeInstances(Element eCoordJob, String eventType, String dataType) throws CoordinatorJobException { 374 Element eventsSpec, dataSpec, instance; 375 List<Element> instanceSpecList; 376 Namespace ns = eCoordJob.getNamespace(); 377 String instanceValue; 378 eventsSpec = eCoordJob.getChild(eventType, ns); 379 if (eventsSpec != null) { 380 dataSpec = eventsSpec.getChild(dataType, ns); 381 if (dataSpec != null) { 382 // In case of input-events, there can be multiple child <instance> datasets. Iterating to ensure none of them have errors 383 instanceSpecList = dataSpec.getChildren("instance", ns); 384 Iterator instanceIter = instanceSpecList.iterator(); 385 while(instanceIter.hasNext()) { 386 instance = ((Element) instanceIter.next()); 387 if(instance.getContentSize() == 0) { //empty string or whitespace 388 throw new CoordinatorJobException(ErrorCode.E1021, "<instance> tag within " + eventType + " is empty!"); 389 } 390 instanceValue = instance.getContent(0).toString(); 391 boolean isInvalid = false; 392 try { 393 isInvalid = evalAction.checkForExistence(instanceValue, ","); 394 } catch (Exception e) { 395 handleELParseException(eventType, dataType, instanceValue); 396 } 397 if (isInvalid) { // reaching this block implies instance is not empty i.e. length > 0 398 handleExpresionWithMultipleInstances(eventType, dataType, instanceValue); 399 } 400 } 401 402 // In case of input-events, there can be multiple child <start-instance> datasets. Iterating to ensure none of them have errors 403 instanceSpecList = dataSpec.getChildren("start-instance", ns); 404 instanceIter = instanceSpecList.iterator(); 405 while(instanceIter.hasNext()) { 406 instance = ((Element) instanceIter.next()); 407 if(instance.getContentSize() == 0) { //empty string or whitespace 408 throw new CoordinatorJobException(ErrorCode.E1021, "<start-instance> tag within " + eventType + " is empty!"); 409 } 410 instanceValue = instance.getContent(0).toString(); 411 boolean isInvalid = false; 412 try { 413 isInvalid = evalAction.checkForExistence(instanceValue, ","); 414 } catch (Exception e) { 415 handleELParseException(eventType, dataType, instanceValue); 416 } 417 if (isInvalid) { // reaching this block implies start instance is not empty i.e. length > 0 418 handleExpresionWithStartMultipleInstances(eventType, dataType, instanceValue); 419 } 420 } 421 422 // In case of input-events, there can be multiple child <end-instance> datasets. Iterating to ensure none of them have errors 423 instanceSpecList = dataSpec.getChildren("end-instance", ns); 424 instanceIter = instanceSpecList.iterator(); 425 while(instanceIter.hasNext()) { 426 instance = ((Element) instanceIter.next()); 427 if(instance.getContentSize() == 0) { //empty string or whitespace 428 throw new CoordinatorJobException(ErrorCode.E1021, "<end-instance> tag within " + eventType + " is empty!"); 429 } 430 instanceValue = instance.getContent(0).toString(); 431 boolean isInvalid = false; 432 try { 433 isInvalid = evalAction.checkForExistence(instanceValue, ","); 434 } catch (Exception e) { 435 handleELParseException(eventType, dataType, instanceValue); 436 } 437 if (isInvalid) { // reaching this block implies instance is not empty i.e. length > 0 438 handleExpresionWithMultipleEndInstances(eventType, dataType, instanceValue); 439 } 440 } 441 442 } 443 } 444 } 445 446 private void handleELParseException(String eventType, String dataType, String instanceValue) 447 throws CoordinatorJobException { 448 String correctAction = null; 449 if(dataType.equals(COORD_INPUT_EVENTS_DATA_IN)) { 450 correctAction = "Coordinator app definition should have valid <instance> tag for data-in"; 451 } else if(dataType.equals(COORD_OUTPUT_EVENTS_DATA_OUT)) { 452 correctAction = "Coordinator app definition should have valid <instance> tag for data-out"; 453 } 454 throw new CoordinatorJobException(ErrorCode.E1021, eventType + " instance '" + instanceValue 455 + "' is not valid. Coordinator job NOT SUBMITTED. " + correctAction); 456 } 457 458 private void handleExpresionWithMultipleInstances(String eventType, String dataType, String instanceValue) 459 throws CoordinatorJobException { 460 String correctAction = null; 461 if(dataType.equals(COORD_INPUT_EVENTS_DATA_IN)) { 462 correctAction = "Coordinator app definition should have separate <instance> tag per data-in instance"; 463 } else if(dataType.equals(COORD_OUTPUT_EVENTS_DATA_OUT)) { 464 correctAction = "Coordinator app definition can have only one <instance> tag per data-out instance"; 465 } 466 throw new CoordinatorJobException(ErrorCode.E1021, eventType + " instance '" + instanceValue 467 + "' contains more than one date instance. Coordinator job NOT SUBMITTED. " + correctAction); 468 } 469 470 private void handleExpresionWithStartMultipleInstances(String eventType, String dataType, String instanceValue) 471 throws CoordinatorJobException { 472 String correctAction = "Coordinator app definition should not have multiple start-instances"; 473 throw new CoordinatorJobException(ErrorCode.E1021, eventType + " start-instance '" + instanceValue 474 + "' contains more than one date start-instance. Coordinator job NOT SUBMITTED. " + correctAction); 475 } 476 477 private void handleExpresionWithMultipleEndInstances(String eventType, String dataType, String instanceValue) 478 throws CoordinatorJobException { 479 String correctAction = "Coordinator app definition should not have multiple end-instances"; 480 throw new CoordinatorJobException(ErrorCode.E1021, eventType + " end-instance '" + instanceValue 481 + "' contains more than one date end-instance. Coordinator job NOT SUBMITTED. " + correctAction); 482 } 483 /** 484 * Read the application XML and validate against coordinator Schema 485 * 486 * @return validated coordinator XML 487 * @throws CoordinatorJobException thrown if unable to read or validate coordinator xml 488 */ 489 protected String readAndValidateXml() throws CoordinatorJobException { 490 String appPath = ParamChecker.notEmpty(conf.get(OozieClient.COORDINATOR_APP_PATH), 491 OozieClient.COORDINATOR_APP_PATH); 492 String coordXml = readDefinition(appPath); 493 validateXml(coordXml); 494 return coordXml; 495 } 496 497 /** 498 * Validate against Coordinator XSD file 499 * 500 * @param xmlContent : Input coordinator xml 501 * @throws CoordinatorJobException thrown if unable to validate coordinator xml 502 */ 503 private void validateXml(String xmlContent) throws CoordinatorJobException { 504 try { 505 Validator validator = Services.get().get(SchemaService.class).getValidator(SchemaName.COORDINATOR); 506 validator.validate(new StreamSource(new StringReader(xmlContent))); 507 } 508 catch (SAXException ex) { 509 LOG.warn("SAXException :", ex); 510 throw new CoordinatorJobException(ErrorCode.E0701, ex.getMessage(), ex); 511 } 512 catch (IOException ex) { 513 LOG.warn("IOException :", ex); 514 throw new CoordinatorJobException(ErrorCode.E0702, ex.getMessage(), ex); 515 } 516 } 517 518 /** 519 * Read the application XML schema namespace 520 * 521 * @param coordXmlElement input coordinator xml Element 522 * @return app xml namespace 523 * @throws CoordinatorJobException 524 */ 525 private String readAppNamespace(Element coordXmlElement) throws CoordinatorJobException { 526 Namespace ns = coordXmlElement.getNamespace(); 527 if (ns != null && bundleId != null && ns.getURI().equals(SchemaService.COORDINATOR_NAMESPACE_URI_1)) { 528 throw new CoordinatorJobException(ErrorCode.E1319, "bundle app can not submit coordinator namespace " 529 + SchemaService.COORDINATOR_NAMESPACE_URI_1 + ", please use 0.2 or later"); 530 } 531 if (ns != null) { 532 return ns.getURI(); 533 } 534 else { 535 throw new CoordinatorJobException(ErrorCode.E0700, "the application xml namespace is not given"); 536 } 537 } 538 539 /** 540 * Merge default configuration with user-defined configuration. 541 * 542 * @throws CommandException thrown if failed to read or merge configurations 543 */ 544 protected void mergeDefaultConfig() throws CommandException { 545 Path configDefault = null; 546 try { 547 String coordAppPathStr = conf.get(OozieClient.COORDINATOR_APP_PATH); 548 Path coordAppPath = new Path(coordAppPathStr); 549 String user = ParamChecker.notEmpty(conf.get(OozieClient.USER_NAME), OozieClient.USER_NAME); 550 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 551 Configuration fsConf = has.createJobConf(coordAppPath.toUri().getAuthority()); 552 FileSystem fs = has.createFileSystem(user, coordAppPath.toUri(), fsConf); 553 554 // app path could be a directory 555 if (!fs.isFile(coordAppPath)) { 556 configDefault = new Path(coordAppPath, CONFIG_DEFAULT); 557 } else { 558 configDefault = new Path(coordAppPath.getParent(), CONFIG_DEFAULT); 559 } 560 561 if (fs.exists(configDefault)) { 562 Configuration defaultConf = new XConfiguration(fs.open(configDefault)); 563 PropertiesUtils.checkDisallowedProperties(defaultConf, DISALLOWED_USER_PROPERTIES); 564 PropertiesUtils.checkDefaultDisallowedProperties(defaultConf); 565 XConfiguration.injectDefaults(defaultConf, conf); 566 } 567 else { 568 LOG.info("configDefault Doesn't exist " + configDefault); 569 } 570 PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_USER_PROPERTIES); 571 572 // Resolving all variables in the job properties. 573 // This ensures the Hadoop Configuration semantics is preserved. 574 XConfiguration resolvedVarsConf = new XConfiguration(); 575 for (Map.Entry<String, String> entry : conf) { 576 resolvedVarsConf.set(entry.getKey(), conf.get(entry.getKey())); 577 } 578 conf = resolvedVarsConf; 579 } 580 catch (IOException e) { 581 throw new CommandException(ErrorCode.E0702, e.getMessage() + " : Problem reading default config " 582 + configDefault, e); 583 } 584 catch (HadoopAccessorException e) { 585 throw new CommandException(e); 586 } 587 LOG.debug("Merged CONF :" + XmlUtils.prettyPrint(conf).toString()); 588 } 589 590 /** 591 * The method resolve all the variables that are defined in configuration. It also include the data set definition 592 * from dataset file into XML. 593 * 594 * @param appXml : Original job XML 595 * @param conf : Configuration of the job 596 * @param coordJob : Coordinator job bean to be populated. 597 * @return Resolved and modified job XML element. 598 * @throws CoordinatorJobException thrown if failed to resolve basic entities or include referred datasets 599 * @throws Exception thrown if failed to resolve basic entities or include referred datasets 600 */ 601 public Element basicResolveAndIncludeDS(String appXml, Configuration conf, CoordinatorJobBean coordJob) 602 throws CoordinatorJobException, Exception { 603 Element basicResolvedApp = resolveInitial(conf, appXml, coordJob); 604 includeDataSets(basicResolvedApp, conf); 605 return basicResolvedApp; 606 } 607 608 /** 609 * Insert data set into data-in and data-out tags. 610 * 611 * @param eAppXml : coordinator application XML 612 * @param eDatasets : DataSet XML 613 */ 614 @SuppressWarnings("unchecked") 615 private void insertDataSet(Element eAppXml, Element eDatasets) { 616 // Adding DS definition in the coordinator XML 617 Element inputList = eAppXml.getChild("input-events", eAppXml.getNamespace()); 618 if (inputList != null) { 619 for (Element dataIn : (List<Element>) inputList.getChildren("data-in", eAppXml.getNamespace())) { 620 Element eDataset = findDataSet(eDatasets, dataIn.getAttributeValue("dataset")); 621 dataIn.getContent().add(0, eDataset); 622 } 623 } 624 Element outputList = eAppXml.getChild("output-events", eAppXml.getNamespace()); 625 if (outputList != null) { 626 for (Element dataOut : (List<Element>) outputList.getChildren("data-out", eAppXml.getNamespace())) { 627 Element eDataset = findDataSet(eDatasets, dataOut.getAttributeValue("dataset")); 628 dataOut.getContent().add(0, eDataset); 629 } 630 } 631 } 632 633 /** 634 * Find a specific dataset from a list of Datasets. 635 * 636 * @param eDatasets : List of data sets 637 * @param name : queried data set name 638 * @return one Dataset element. otherwise throw Exception 639 */ 640 @SuppressWarnings("unchecked") 641 private static Element findDataSet(Element eDatasets, String name) { 642 for (Element eDataset : (List<Element>) eDatasets.getChildren("dataset", eDatasets.getNamespace())) { 643 if (eDataset.getAttributeValue("name").equals(name)) { 644 eDataset = (Element) eDataset.clone(); 645 eDataset.detach(); 646 return eDataset; 647 } 648 } 649 throw new RuntimeException("undefined dataset: " + name); 650 } 651 652 /** 653 * Initialize all the required EL Evaluators. 654 */ 655 protected void initEvaluators() { 656 evalFreq = CoordELEvaluator.createELEvaluatorForGroup(conf, "coord-job-submit-freq"); 657 evalNofuncs = CoordELEvaluator.createELEvaluatorForGroup(conf, "coord-job-submit-nofuncs"); 658 evalInst = CoordELEvaluator.createELEvaluatorForGroup(conf, "coord-job-submit-instances"); 659 evalAction = CoordELEvaluator.createELEvaluatorForGroup(conf, "coord-action-start"); 660 evalTimeout = CoordELEvaluator.createELEvaluatorForGroup(conf, "coord-job-wait-timeout"); 661 evalInitialInstance = CoordELEvaluator.createELEvaluatorForGroup(conf, "coord-job-submit-initial-instance"); 662 663 } 664 665 /** 666 * Resolve basic entities using job Configuration. 667 * 668 * @param conf :Job configuration 669 * @param appXml : Original job XML 670 * @param coordJob : Coordinator job bean to be populated. 671 * @return Resolved job XML element. 672 * @throws CoordinatorJobException thrown if failed to resolve basic entities 673 * @throws Exception thrown if failed to resolve basic entities 674 */ 675 @SuppressWarnings("unchecked") 676 protected Element resolveInitial(Configuration conf, String appXml, CoordinatorJobBean coordJob) 677 throws CoordinatorJobException, Exception { 678 Element eAppXml = XmlUtils.parseXml(appXml); 679 // job's main attributes 680 // frequency 681 String val = resolveAttribute("frequency", eAppXml, evalFreq); 682 int ival = 0; 683 684 val = ParamChecker.checkFrequency(val); 685 coordJob.setFrequency(val); 686 TimeUnit tmp = (evalFreq.getVariable("timeunit") == null) ? TimeUnit.MINUTE : ((TimeUnit) evalFreq 687 .getVariable("timeunit")); 688 try { 689 Integer.parseInt(val); 690 } 691 catch (NumberFormatException ex) { 692 tmp=TimeUnit.CRON; 693 } 694 695 addAnAttribute("freq_timeunit", eAppXml, tmp.toString()); 696 // TimeUnit 697 coordJob.setTimeUnit(CoordinatorJob.Timeunit.valueOf(tmp.toString())); 698 // End Of Duration 699 tmp = evalFreq.getVariable("endOfDuration") == null ? TimeUnit.NONE : ((TimeUnit) evalFreq 700 .getVariable("endOfDuration")); 701 addAnAttribute("end_of_duration", eAppXml, tmp.toString()); 702 // coordJob.setEndOfDuration(tmp) // TODO: Add new attribute in Job bean 703 704 // Application name 705 if (this.coordName == null) { 706 String name = ELUtils.resolveAppName(eAppXml.getAttribute("name").getValue(), conf); 707 coordJob.setAppName(name); 708 } 709 else { 710 // this coord job is created from bundle 711 coordJob.setAppName(this.coordName); 712 } 713 714 // start time 715 val = resolveAttribute("start", eAppXml, evalNofuncs); 716 ParamChecker.checkDateOozieTZ(val, "start"); 717 coordJob.setStartTime(DateUtils.parseDateOozieTZ(val)); 718 // end time 719 val = resolveAttribute("end", eAppXml, evalNofuncs); 720 ParamChecker.checkDateOozieTZ(val, "end"); 721 coordJob.setEndTime(DateUtils.parseDateOozieTZ(val)); 722 // Time zone 723 val = resolveAttribute("timezone", eAppXml, evalNofuncs); 724 ParamChecker.checkTimeZone(val, "timezone"); 725 coordJob.setTimeZone(val); 726 727 // controls 728 val = resolveTagContents("timeout", eAppXml.getChild("controls", eAppXml.getNamespace()), evalTimeout); 729 if (val != null && val != "") { 730 int t = Integer.parseInt(val); 731 tmp = (evalTimeout.getVariable("timeunit") == null) ? TimeUnit.MINUTE : ((TimeUnit) evalTimeout 732 .getVariable("timeunit")); 733 switch (tmp) { 734 case HOUR: 735 val = String.valueOf(t * 60); 736 break; 737 case DAY: 738 val = String.valueOf(t * 60 * 24); 739 break; 740 case MONTH: 741 val = String.valueOf(t * 60 * 24 * 30); 742 break; 743 default: 744 break; 745 } 746 } 747 else { 748 val = ConfigurationService.get(CONF_DEFAULT_TIMEOUT_NORMAL); 749 } 750 751 ival = ParamChecker.checkInteger(val, "timeout"); 752 if (ival < 0 || ival > ConfigurationService.getInt(CONF_DEFAULT_MAX_TIMEOUT)) { 753 ival = ConfigurationService.getInt(CONF_DEFAULT_MAX_TIMEOUT); 754 } 755 coordJob.setTimeout(ival); 756 757 val = resolveTagContents("concurrency", eAppXml.getChild("controls", eAppXml.getNamespace()), evalNofuncs); 758 if (val == null || val.isEmpty()) { 759 val = ConfigurationService.get(CONF_DEFAULT_CONCURRENCY); 760 } 761 ival = ParamChecker.checkInteger(val, "concurrency"); 762 coordJob.setConcurrency(ival); 763 764 val = resolveTagContents("throttle", eAppXml.getChild("controls", eAppXml.getNamespace()), evalNofuncs); 765 if (val == null || val.isEmpty()) { 766 int defaultThrottle = ConfigurationService.getInt(CONF_DEFAULT_THROTTLE); 767 ival = defaultThrottle; 768 } 769 else { 770 ival = ParamChecker.checkInteger(val, "throttle"); 771 } 772 int maxQueue = ConfigurationService.getInt(CONF_QUEUE_SIZE); 773 float factor = ConfigurationService.getFloat(CONF_MAT_THROTTLING_FACTOR); 774 int maxThrottle = (int) (maxQueue * factor); 775 if (ival > maxThrottle || ival < 1) { 776 ival = maxThrottle; 777 } 778 LOG.debug("max throttle " + ival); 779 coordJob.setMatThrottling(ival); 780 781 val = resolveTagContents("execution", eAppXml.getChild("controls", eAppXml.getNamespace()), evalNofuncs); 782 if (val == "") { 783 val = Execution.FIFO.toString(); 784 } 785 coordJob.setExecutionOrder(Execution.valueOf(val)); 786 String[] acceptedVals = { Execution.LIFO.toString(), Execution.FIFO.toString(), Execution.LAST_ONLY.toString(), 787 Execution.NONE.toString()}; 788 ParamChecker.isMember(val, acceptedVals, "execution"); 789 790 // datasets 791 resolveTagContents("include", eAppXml.getChild("datasets", eAppXml.getNamespace()), evalNofuncs); 792 // for each data set 793 resolveDataSets(eAppXml); 794 HashMap<String, String> dataNameList = new HashMap<String, String>(); 795 resolveIODataset(eAppXml); 796 resolveIOEvents(eAppXml, dataNameList); 797 798 resolveTagContents("app-path", eAppXml.getChild("action", eAppXml.getNamespace()).getChild("workflow", 799 eAppXml.getNamespace()), evalNofuncs); 800 // TODO: If action or workflow tag is missing, NullPointerException will 801 // occur 802 Element configElem = eAppXml.getChild("action", eAppXml.getNamespace()).getChild("workflow", 803 eAppXml.getNamespace()).getChild("configuration", eAppXml.getNamespace()); 804 evalData = CoordELEvaluator.createELEvaluatorForDataEcho(conf, "coord-job-submit-data", dataNameList); 805 if (configElem != null) { 806 for (Element propElem : (List<Element>) configElem.getChildren("property", configElem.getNamespace())) { 807 resolveTagContents("name", propElem, evalData); 808 // Want to check the data-integrity but don't want to modify the 809 // XML 810 // for properties only 811 Element tmpProp = (Element) propElem.clone(); 812 resolveTagContents("value", tmpProp, evalData); 813 } 814 } 815 evalSla = CoordELEvaluator.createELEvaluatorForDataAndConf(conf, "coord-sla-submit", dataNameList); 816 resolveSLA(eAppXml, coordJob); 817 return eAppXml; 818 } 819 820 /** 821 * Resolve SLA events 822 * 823 * @param eAppXml job XML 824 * @param coordJob coordinator job bean 825 * @throws CommandException thrown if failed to resolve sla events 826 */ 827 private void resolveSLA(Element eAppXml, CoordinatorJobBean coordJob) throws CommandException { 828 Element eSla = XmlUtils.getSLAElement(eAppXml.getChild("action", eAppXml.getNamespace())); 829 830 if (eSla != null) { 831 String slaXml = XmlUtils.prettyPrint(eSla).toString(); 832 try { 833 // EL evaluation 834 slaXml = evalSla.evaluate(slaXml, String.class); 835 // Validate against semantic SXD 836 XmlUtils.validateData(slaXml, SchemaName.SLA_ORIGINAL); 837 } 838 catch (Exception e) { 839 throw new CommandException(ErrorCode.E1004, "Validation ERROR :" + e.getMessage(), e); 840 } 841 } 842 } 843 844 /** 845 * Resolve input-events/data-in and output-events/data-out tags. 846 * 847 * @param eJobOrg : Job element 848 * @throws CoordinatorJobException thrown if failed to resolve input and output events 849 */ 850 @SuppressWarnings("unchecked") 851 private void resolveIOEvents(Element eJobOrg, HashMap<String, String> dataNameList) throws CoordinatorJobException { 852 // Resolving input-events/data-in 853 // Clone the job and don't update anything in the original 854 Element eJob = (Element) eJobOrg.clone(); 855 Element inputList = eJob.getChild("input-events", eJob.getNamespace()); 856 if (inputList != null) { 857 TreeSet<String> eventNameSet = new TreeSet<String>(); 858 for (Element dataIn : (List<Element>) inputList.getChildren("data-in", eJob.getNamespace())) { 859 String dataInName = dataIn.getAttributeValue("name"); 860 dataNameList.put(dataInName, "data-in"); 861 // check whether there is any duplicate data-in name 862 if (eventNameSet.contains(dataInName)) { 863 throw new RuntimeException("Duplicate dataIn name " + dataInName); 864 } 865 else { 866 eventNameSet.add(dataInName); 867 } 868 resolveTagContents("instance", dataIn, evalInst); 869 resolveTagContents("start-instance", dataIn, evalInst); 870 resolveTagContents("end-instance", dataIn, evalInst); 871 872 } 873 } 874 // Resolving output-events/data-out 875 Element outputList = eJob.getChild("output-events", eJob.getNamespace()); 876 if (outputList != null) { 877 TreeSet<String> eventNameSet = new TreeSet<String>(); 878 for (Element dataOut : (List<Element>) outputList.getChildren("data-out", eJob.getNamespace())) { 879 String dataOutName = dataOut.getAttributeValue("name"); 880 dataNameList.put(dataOutName, "data-out"); 881 // check whether there is any duplicate data-out name 882 if (eventNameSet.contains(dataOutName)) { 883 throw new RuntimeException("Duplicate dataIn name " + dataOutName); 884 } 885 else { 886 eventNameSet.add(dataOutName); 887 } 888 resolveTagContents("instance", dataOut, evalInst); 889 890 } 891 } 892 893 } 894 895 /** 896 * Resolve input-events/dataset and output-events/dataset tags. 897 * 898 * @param eJob : Job element 899 * @throws CoordinatorJobException thrown if failed to resolve input and output events 900 */ 901 @SuppressWarnings("unchecked") 902 private void resolveIODataset(Element eAppXml) throws CoordinatorJobException { 903 // Resolving input-events/data-in 904 Element inputList = eAppXml.getChild("input-events", eAppXml.getNamespace()); 905 if (inputList != null) { 906 for (Element dataIn : (List<Element>) inputList.getChildren("data-in", eAppXml.getNamespace())) { 907 resolveAttribute("dataset", dataIn, evalInst); 908 909 } 910 } 911 // Resolving output-events/data-out 912 Element outputList = eAppXml.getChild("output-events", eAppXml.getNamespace()); 913 if (outputList != null) { 914 for (Element dataOut : (List<Element>) outputList.getChildren("data-out", eAppXml.getNamespace())) { 915 resolveAttribute("dataset", dataOut, evalInst); 916 917 } 918 } 919 920 } 921 922 923 /** 924 * Add an attribute into XML element. 925 * 926 * @param attrName :attribute name 927 * @param elem : Element to add attribute 928 * @param value :Value of attribute 929 */ 930 private void addAnAttribute(String attrName, Element elem, String value) { 931 elem.setAttribute(attrName, value); 932 } 933 934 /** 935 * Resolve datasets using job configuration. 936 * 937 * @param eAppXml : Job Element XML 938 * @throws Exception thrown if failed to resolve datasets 939 */ 940 @SuppressWarnings("unchecked") 941 private void resolveDataSets(Element eAppXml) throws Exception { 942 Element datasetList = eAppXml.getChild("datasets", eAppXml.getNamespace()); 943 if (datasetList != null) { 944 945 List<Element> dsElems = datasetList.getChildren("dataset", eAppXml.getNamespace()); 946 resolveDataSets(dsElems); 947 resolveTagContents("app-path", eAppXml.getChild("action", eAppXml.getNamespace()).getChild("workflow", 948 eAppXml.getNamespace()), evalNofuncs); 949 } 950 } 951 952 /** 953 * Resolve datasets using job configuration. 954 * 955 * @param dsElems : Data set XML element. 956 * @throws CoordinatorJobException thrown if failed to resolve datasets 957 */ 958 private void resolveDataSets(List<Element> dsElems) throws CoordinatorJobException { 959 for (Element dsElem : dsElems) { 960 // Setting up default TimeUnit and EndOFDuraion 961 evalFreq.setVariable("timeunit", TimeUnit.MINUTE); 962 evalFreq.setVariable("endOfDuration", TimeUnit.NONE); 963 964 String val = resolveAttribute("frequency", dsElem, evalFreq); 965 int ival = ParamChecker.checkInteger(val, "frequency"); 966 ParamChecker.checkGTZero(ival, "frequency"); 967 addAnAttribute("freq_timeunit", dsElem, evalFreq.getVariable("timeunit") == null ? TimeUnit.MINUTE 968 .toString() : ((TimeUnit) evalFreq.getVariable("timeunit")).toString()); 969 addAnAttribute("end_of_duration", dsElem, evalFreq.getVariable("endOfDuration") == null ? TimeUnit.NONE 970 .toString() : ((TimeUnit) evalFreq.getVariable("endOfDuration")).toString()); 971 val = resolveAttribute("initial-instance", dsElem, evalInitialInstance); 972 ParamChecker.checkDateOozieTZ(val, "initial-instance"); 973 checkInitialInstance(val); 974 val = resolveAttribute("timezone", dsElem, evalNofuncs); 975 ParamChecker.checkTimeZone(val, "timezone"); 976 resolveTagContents("uri-template", dsElem, evalNofuncs); 977 resolveTagContents("done-flag", dsElem, evalNofuncs); 978 } 979 } 980 981 /** 982 * Resolve the content of a tag. 983 * 984 * @param tagName : Tag name of job XML i.e. <timeout> 10 </timeout> 985 * @param elem : Element where the tag exists. 986 * @param eval : EL evealuator 987 * @return Resolved tag content. 988 * @throws CoordinatorJobException thrown if failed to resolve tag content 989 */ 990 @SuppressWarnings("unchecked") 991 private String resolveTagContents(String tagName, Element elem, ELEvaluator eval) throws CoordinatorJobException { 992 String ret = ""; 993 if (elem != null) { 994 for (Element tagElem : (List<Element>) elem.getChildren(tagName, elem.getNamespace())) { 995 if (tagElem != null) { 996 String updated; 997 try { 998 updated = CoordELFunctions.evalAndWrap(eval, tagElem.getText().trim()); 999 1000 } 1001 catch (Exception e) { 1002 throw new CoordinatorJobException(ErrorCode.E1004, e.getMessage(), e); 1003 } 1004 tagElem.removeContent(); 1005 tagElem.addContent(updated); 1006 ret += updated; 1007 } 1008 } 1009 } 1010 return ret; 1011 } 1012 1013 /** 1014 * Resolve an attribute value. 1015 * 1016 * @param attrName : Attribute name. 1017 * @param elem : XML Element where attribute is defiend 1018 * @param eval : ELEvaluator used to resolve 1019 * @return Resolved attribute value 1020 * @throws CoordinatorJobException thrown if failed to resolve an attribute value 1021 */ 1022 private String resolveAttribute(String attrName, Element elem, ELEvaluator eval) throws CoordinatorJobException { 1023 Attribute attr = elem.getAttribute(attrName); 1024 String val = null; 1025 if (attr != null) { 1026 try { 1027 val = CoordELFunctions.evalAndWrap(eval, attr.getValue().trim()); 1028 } 1029 catch (Exception e) { 1030 throw new CoordinatorJobException(ErrorCode.E1004, e.getMessage(), e); 1031 } 1032 attr.setValue(val); 1033 } 1034 return val; 1035 } 1036 1037 /** 1038 * Include referred datasets into XML. 1039 * 1040 * @param resolvedXml : Job XML element. 1041 * @param conf : Job configuration 1042 * @throws CoordinatorJobException thrown if failed to include referred datasets into XML 1043 */ 1044 @SuppressWarnings("unchecked") 1045 protected void includeDataSets(Element resolvedXml, Configuration conf) throws CoordinatorJobException { 1046 Element datasets = resolvedXml.getChild("datasets", resolvedXml.getNamespace()); 1047 Element allDataSets = new Element("all_datasets", resolvedXml.getNamespace()); 1048 List<String> dsList = new ArrayList<String>(); 1049 if (datasets != null) { 1050 for (Element includeElem : (List<Element>) datasets.getChildren("include", datasets.getNamespace())) { 1051 String incDSFile = includeElem.getTextTrim(); 1052 includeOneDSFile(incDSFile, dsList, allDataSets, datasets.getNamespace()); 1053 } 1054 for (Element e : (List<Element>) datasets.getChildren("dataset", datasets.getNamespace())) { 1055 String dsName = e.getAttributeValue("name"); 1056 if (dsList.contains(dsName)) {// Override with this DS 1057 // Remove duplicate 1058 removeDataSet(allDataSets, dsName); 1059 } 1060 else { 1061 dsList.add(dsName); 1062 } 1063 allDataSets.addContent((Element) e.clone()); 1064 } 1065 } 1066 insertDataSet(resolvedXml, allDataSets); 1067 resolvedXml.removeChild("datasets", resolvedXml.getNamespace()); 1068 } 1069 1070 /** 1071 * Include one dataset file. 1072 * 1073 * @param incDSFile : Include data set filename. 1074 * @param dsList :List of dataset names to verify the duplicate. 1075 * @param allDataSets : Element that includes all dataset definitions. 1076 * @param dsNameSpace : Data set name space 1077 * @throws CoordinatorJobException thrown if failed to include one dataset file 1078 */ 1079 @SuppressWarnings("unchecked") 1080 private void includeOneDSFile(String incDSFile, List<String> dsList, Element allDataSets, Namespace dsNameSpace) 1081 throws CoordinatorJobException { 1082 Element tmpDataSets = null; 1083 try { 1084 String dsXml = readDefinition(incDSFile); 1085 LOG.debug("DSFILE :" + incDSFile + "\n" + dsXml); 1086 tmpDataSets = XmlUtils.parseXml(dsXml); 1087 } 1088 catch (JDOMException e) { 1089 LOG.warn("Error parsing included dataset [{0}]. Message [{1}]", incDSFile, e.getMessage()); 1090 throw new CoordinatorJobException(ErrorCode.E0700, e.getMessage()); 1091 } 1092 resolveDataSets(tmpDataSets.getChildren("dataset")); 1093 for (Element e : (List<Element>) tmpDataSets.getChildren("dataset")) { 1094 String dsName = e.getAttributeValue("name"); 1095 if (dsList.contains(dsName)) { 1096 throw new RuntimeException("Duplicate Dataset " + dsName); 1097 } 1098 dsList.add(dsName); 1099 Element tmp = (Element) e.clone(); 1100 // TODO: Don't like to over-write the external/include DS's namespace 1101 tmp.setNamespace(dsNameSpace); 1102 tmp.getChild("uri-template").setNamespace(dsNameSpace); 1103 if (e.getChild("done-flag") != null) { 1104 tmp.getChild("done-flag").setNamespace(dsNameSpace); 1105 } 1106 allDataSets.addContent(tmp); 1107 } 1108 // nested include 1109 for (Element includeElem : (List<Element>) tmpDataSets.getChildren("include", tmpDataSets.getNamespace())) { 1110 String incFile = includeElem.getTextTrim(); 1111 includeOneDSFile(incFile, dsList, allDataSets, dsNameSpace); 1112 } 1113 } 1114 1115 /** 1116 * Remove a dataset from a list of dataset. 1117 * 1118 * @param eDatasets : List of dataset 1119 * @param name : Dataset name to be removed. 1120 */ 1121 @SuppressWarnings("unchecked") 1122 private static void removeDataSet(Element eDatasets, String name) { 1123 for (Element eDataset : (List<Element>) eDatasets.getChildren("dataset", eDatasets.getNamespace())) { 1124 if (eDataset.getAttributeValue("name").equals(name)) { 1125 eDataset.detach(); 1126 return; 1127 } 1128 } 1129 throw new RuntimeException("undefined dataset: " + name); 1130 } 1131 1132 /** 1133 * Read coordinator definition. 1134 * 1135 * @param appPath application path. 1136 * @return coordinator definition. 1137 * @throws CoordinatorJobException thrown if the definition could not be read. 1138 */ 1139 protected String readDefinition(String appPath) throws CoordinatorJobException { 1140 String user = ParamChecker.notEmpty(conf.get(OozieClient.USER_NAME), OozieClient.USER_NAME); 1141 // Configuration confHadoop = CoordUtils.getHadoopConf(conf); 1142 try { 1143 URI uri = new URI(appPath); 1144 LOG.debug("user =" + user); 1145 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 1146 Configuration fsConf = has.createJobConf(uri.getAuthority()); 1147 FileSystem fs = has.createFileSystem(user, uri, fsConf); 1148 Path appDefPath = null; 1149 1150 // app path could be a directory 1151 Path path = new Path(uri.getPath()); 1152 // check file exists for dataset include file, app xml already checked 1153 if (!fs.exists(path)) { 1154 throw new URISyntaxException(path.toString(), "path not existed : " + path.toString()); 1155 } 1156 if (!fs.isFile(path)) { 1157 appDefPath = new Path(path, COORDINATOR_XML_FILE); 1158 } else { 1159 appDefPath = path; 1160 } 1161 1162 Reader reader = new InputStreamReader(fs.open(appDefPath)); 1163 StringWriter writer = new StringWriter(); 1164 IOUtils.copyCharStream(reader, writer); 1165 return writer.toString(); 1166 } 1167 catch (IOException ex) { 1168 LOG.warn("IOException :" + XmlUtils.prettyPrint(conf), ex); 1169 throw new CoordinatorJobException(ErrorCode.E1001, ex.getMessage(), ex); 1170 } 1171 catch (URISyntaxException ex) { 1172 LOG.warn("URISyException :" + ex.getMessage()); 1173 throw new CoordinatorJobException(ErrorCode.E1002, appPath, ex.getMessage(), ex); 1174 } 1175 catch (HadoopAccessorException ex) { 1176 throw new CoordinatorJobException(ex); 1177 } 1178 catch (Exception ex) { 1179 LOG.warn("Exception :", ex); 1180 throw new CoordinatorJobException(ErrorCode.E1001, ex.getMessage(), ex); 1181 } 1182 } 1183 1184 /** 1185 * Write a coordinator job into database 1186 * 1187 *@param appXML : Coordinator definition xml 1188 * @param eJob : XML element of job 1189 * @param coordJob : Coordinator job bean 1190 * @return Job id 1191 * @throws CommandException thrown if unable to save coordinator job to db 1192 */ 1193 protected String storeToDB(String appXML, Element eJob, CoordinatorJobBean coordJob) throws CommandException { 1194 String jobId = Services.get().get(UUIDService.class).generateId(ApplicationType.COORDINATOR); 1195 coordJob.setId(jobId); 1196 1197 coordJob.setAppPath(conf.get(OozieClient.COORDINATOR_APP_PATH)); 1198 coordJob.setCreatedTime(new Date()); 1199 coordJob.setUser(conf.get(OozieClient.USER_NAME)); 1200 String group = ConfigUtils.getWithDeprecatedCheck(conf, OozieClient.JOB_ACL, OozieClient.GROUP_NAME, null); 1201 coordJob.setGroup(group); 1202 coordJob.setConf(XmlUtils.prettyPrint(conf).toString()); 1203 coordJob.setJobXml(XmlUtils.prettyPrint(eJob).toString()); 1204 coordJob.setLastActionNumber(0); 1205 coordJob.setLastModifiedTime(new Date()); 1206 1207 if (!dryrun) { 1208 coordJob.setLastModifiedTime(new Date()); 1209 try { 1210 CoordJobQueryExecutor.getInstance().insert(coordJob); 1211 } 1212 catch (JPAExecutorException jpaee) { 1213 coordJob.setId(null); 1214 coordJob.setStatus(CoordinatorJob.Status.FAILED); 1215 throw new CommandException(jpaee); 1216 } 1217 } 1218 return jobId; 1219 } 1220 1221 /* 1222 * this method checks if the initial-instance specified for a particular 1223 is not a date earlier than the oozie server default Jan 01, 1970 00:00Z UTC 1224 */ 1225 private void checkInitialInstance(String val) throws CoordinatorJobException, IllegalArgumentException { 1226 Date initialInstance, givenInstance; 1227 try { 1228 initialInstance = DateUtils.parseDateUTC("1970-01-01T00:00Z"); 1229 givenInstance = DateUtils.parseDateOozieTZ(val); 1230 } 1231 catch (Exception e) { 1232 throw new IllegalArgumentException("Unable to parse dataset initial-instance string '" + val + 1233 "' to Date object. ",e); 1234 } 1235 if(givenInstance.compareTo(initialInstance) < 0) { 1236 throw new CoordinatorJobException(ErrorCode.E1021, "Dataset initial-instance " + val + 1237 " is earlier than the default initial instance " + DateUtils.formatDateOozieTZ(initialInstance)); 1238 } 1239 } 1240 1241 /* (non-Javadoc) 1242 * @see org.apache.oozie.command.XCommand#getEntityKey() 1243 */ 1244 @Override 1245 public String getEntityKey() { 1246 return null; 1247 } 1248 1249 /* (non-Javadoc) 1250 * @see org.apache.oozie.command.XCommand#isLockRequired() 1251 */ 1252 @Override 1253 protected boolean isLockRequired() { 1254 return false; 1255 } 1256 1257 /* (non-Javadoc) 1258 * @see org.apache.oozie.command.XCommand#loadState() 1259 */ 1260 @Override 1261 protected void loadState() throws CommandException { 1262 jpaService = Services.get().get(JPAService.class); 1263 if (jpaService == null) { 1264 throw new CommandException(ErrorCode.E0610); 1265 } 1266 coordJob = new CoordinatorJobBean(); 1267 if (this.bundleId != null) { 1268 // this coord job is created from bundle 1269 coordJob.setBundleId(this.bundleId); 1270 // first use bundle id if submit thru bundle 1271 LogUtils.setLogInfo(this.bundleId); 1272 } 1273 if (this.coordName != null) { 1274 // this coord job is created from bundle 1275 coordJob.setAppName(this.coordName); 1276 } 1277 setJob(coordJob); 1278 1279 } 1280 1281 /* (non-Javadoc) 1282 * @see org.apache.oozie.command.XCommand#verifyPrecondition() 1283 */ 1284 @Override 1285 protected void verifyPrecondition() throws CommandException { 1286 1287 } 1288 1289 /* (non-Javadoc) 1290 * @see org.apache.oozie.command.TransitionXCommand#notifyParent() 1291 */ 1292 @Override 1293 public void notifyParent() throws CommandException { 1294 // update bundle action 1295 if (coordJob.getBundleId() != null) { 1296 LOG.debug("Updating bundle record: " + coordJob.getBundleId() + " for coord id: " + coordJob.getId()); 1297 BundleStatusUpdateXCommand bundleStatusUpdate = new BundleStatusUpdateXCommand(coordJob, prevStatus); 1298 bundleStatusUpdate.call(); 1299 } 1300 } 1301 1302 /* (non-Javadoc) 1303 * @see org.apache.oozie.command.TransitionXCommand#updateJob() 1304 */ 1305 @Override 1306 public void updateJob() throws CommandException { 1307 } 1308 1309 /* (non-Javadoc) 1310 * @see org.apache.oozie.command.TransitionXCommand#getJob() 1311 */ 1312 @Override 1313 public Job getJob() { 1314 return coordJob; 1315 } 1316 1317 @Override 1318 public void performWrites() throws CommandException { 1319 } 1320}