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.bundle; 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.Date; 029import java.util.HashSet; 030import java.util.List; 031import java.util.Map; 032import java.util.Set; 033 034import javax.xml.transform.stream.StreamSource; 035import javax.xml.validation.Validator; 036 037import org.apache.hadoop.conf.Configuration; 038import org.apache.hadoop.fs.FileSystem; 039import org.apache.hadoop.fs.Path; 040import org.apache.oozie.BundleJobBean; 041import org.apache.oozie.ErrorCode; 042import org.apache.oozie.client.Job; 043import org.apache.oozie.client.OozieClient; 044import org.apache.oozie.command.CommandException; 045import org.apache.oozie.command.PreconditionException; 046import org.apache.oozie.command.SubmitTransitionXCommand; 047import org.apache.oozie.executor.jpa.BundleJobQueryExecutor; 048import org.apache.oozie.service.ELService; 049import org.apache.oozie.service.HadoopAccessorException; 050import org.apache.oozie.service.HadoopAccessorService; 051import org.apache.oozie.service.SchemaService; 052import org.apache.oozie.service.SchemaService.SchemaName; 053import org.apache.oozie.service.Services; 054import org.apache.oozie.service.UUIDService; 055import org.apache.oozie.service.UUIDService.ApplicationType; 056import org.apache.oozie.util.ConfigUtils; 057import org.apache.oozie.util.DateUtils; 058import org.apache.oozie.util.ELEvaluator; 059import org.apache.oozie.util.ELUtils; 060import org.apache.oozie.util.IOUtils; 061import org.apache.oozie.util.InstrumentUtils; 062import org.apache.oozie.util.LogUtils; 063import org.apache.oozie.util.ParamChecker; 064import org.apache.oozie.util.ParameterVerifier; 065import org.apache.oozie.util.PropertiesUtils; 066import org.apache.oozie.util.XConfiguration; 067import org.apache.oozie.util.XmlUtils; 068import org.jdom.Attribute; 069import org.jdom.Element; 070import org.jdom.JDOMException; 071import org.xml.sax.SAXException; 072 073/** 074 * This Command will submit the bundle. 075 */ 076public class BundleSubmitXCommand extends SubmitTransitionXCommand { 077 078 private Configuration conf; 079 public static final String CONFIG_DEFAULT = "bundle-config-default.xml"; 080 public static final String BUNDLE_XML_FILE = "bundle.xml"; 081 private final BundleJobBean bundleBean = new BundleJobBean(); 082 private String jobId; 083 private static final Set<String> DISALLOWED_USER_PROPERTIES = new HashSet<String>(); 084 private static final Set<String> DISALLOWED_DEFAULT_PROPERTIES = new HashSet<String>(); 085 086 static { 087 String[] badUserProps = { PropertiesUtils.YEAR, PropertiesUtils.MONTH, PropertiesUtils.DAY, 088 PropertiesUtils.HOUR, PropertiesUtils.MINUTE, PropertiesUtils.DAYS, PropertiesUtils.HOURS, 089 PropertiesUtils.MINUTES, PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, 090 PropertiesUtils.TB, PropertiesUtils.PB, PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, 091 PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN, PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS }; 092 PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_USER_PROPERTIES); 093 } 094 095 /** 096 * Constructor to create the bundle submit command. 097 * 098 * @param conf configuration for bundle job 099 */ 100 public BundleSubmitXCommand(Configuration conf) { 101 super("bundle_submit", "bundle_submit", 1); 102 this.conf = ParamChecker.notNull(conf, "conf"); 103 } 104 105 /** 106 * Constructor to create the bundle submit command. 107 * 108 * @param dryrun true if dryrun is enable 109 * @param conf configuration for bundle job 110 */ 111 public BundleSubmitXCommand(boolean dryrun, Configuration conf) { 112 this(conf); 113 this.dryrun = dryrun; 114 } 115 116 /* (non-Javadoc) 117 * @see org.apache.oozie.command.SubmitTransitionXCommand#submit() 118 */ 119 @Override 120 protected String submit() throws CommandException { 121 LOG.info("STARTED Bundle Submit"); 122 try { 123 InstrumentUtils.incrJobCounter(getName(), 1, getInstrumentation()); 124 125 ParameterVerifier.verifyParameters(conf, XmlUtils.parseXml(bundleBean.getOrigJobXml())); 126 127 String jobXmlWithNoComment = XmlUtils.removeComments(this.bundleBean.getOrigJobXml().toString()); 128 // Resolving all variables in the job properties. 129 // This ensures the Hadoop Configuration semantics is preserved. 130 XConfiguration resolvedVarsConf = new XConfiguration(); 131 for (Map.Entry<String, String> entry : conf) { 132 resolvedVarsConf.set(entry.getKey(), conf.get(entry.getKey())); 133 } 134 conf = resolvedVarsConf; 135 136 String resolvedJobXml = resolvedVarsandFunctions(jobXmlWithNoComment, conf); 137 138 //verify the uniqueness of coord names 139 verifyCoordNameUnique(resolvedJobXml); 140 this.jobId = storeToDB(bundleBean, resolvedJobXml); 141 LogUtils.setLogInfo(bundleBean); 142 143 if (dryrun) { 144 Date startTime = bundleBean.getStartTime(); 145 long startTimeMilli = startTime.getTime(); 146 long endTimeMilli = startTimeMilli + (3600 * 1000); 147 Date jobEndTime = bundleBean.getEndTime(); 148 Date endTime = new Date(endTimeMilli); 149 if (endTime.compareTo(jobEndTime) > 0) { 150 endTime = jobEndTime; 151 } 152 jobId = bundleBean.getId(); 153 LOG.info("[" + jobId + "]: Update status to PREP"); 154 bundleBean.setStatus(Job.Status.PREP); 155 try { 156 new XConfiguration(new StringReader(bundleBean.getConf())); 157 } 158 catch (IOException e1) { 159 LOG.warn("Configuration parse error. read from DB :" + bundleBean.getConf(), e1); 160 } 161 String output = bundleBean.getJobXml() + System.getProperty("line.separator"); 162 return output; 163 } 164 else { 165 if (bundleBean.getKickoffTime() == null) { 166 // If there is no KickOffTime, default kickoff is NOW. 167 LOG.debug("Since kickoff time is not defined for job id " + jobId 168 + ". Queuing and BundleStartXCommand immediately after submission"); 169 queue(new BundleStartXCommand(jobId)); 170 } 171 } 172 } 173 catch (Exception ex) { 174 throw new CommandException(ErrorCode.E1310, ex.getMessage(), ex); 175 } 176 LOG.info("ENDED Bundle Submit"); 177 return this.jobId; 178 } 179 180 /* (non-Javadoc) 181 * @see org.apache.oozie.command.TransitionXCommand#notifyParent() 182 */ 183 @Override 184 public void notifyParent() throws CommandException { 185 } 186 187 /* (non-Javadoc) 188 * @see org.apache.oozie.command.XCommand#getEntityKey() 189 */ 190 @Override 191 public String getEntityKey() { 192 return null; 193 } 194 195 /* (non-Javadoc) 196 * @see org.apache.oozie.command.XCommand#isLockRequired() 197 */ 198 @Override 199 protected boolean isLockRequired() { 200 return false; 201 } 202 203 @Override 204 protected void loadState() throws CommandException { 205 } 206 207 @Override 208 protected void verifyPrecondition() throws CommandException, PreconditionException { 209 } 210 211 @Override 212 protected void eagerLoadState() throws CommandException { 213 } 214 215 @Override 216 protected void eagerVerifyPrecondition() throws CommandException, PreconditionException { 217 try { 218 mergeDefaultConfig(); 219 String appXml = readAndValidateXml(); 220 bundleBean.setOrigJobXml(appXml); 221 LOG.debug("jobXml after initial validation " + XmlUtils.prettyPrint(appXml).toString()); 222 } 223 catch (BundleJobException ex) { 224 LOG.warn("BundleJobException: ", ex); 225 throw new CommandException(ex); 226 } 227 catch (IllegalArgumentException iex) { 228 LOG.warn("IllegalArgumentException: ", iex); 229 throw new CommandException(ErrorCode.E1310, iex.getMessage(), iex); 230 } 231 catch (Exception ex) { 232 LOG.warn("Exception: ", ex); 233 throw new CommandException(ErrorCode.E1310, ex.getMessage(), ex); 234 } 235 } 236 237 /** 238 * Merge default configuration with user-defined configuration. 239 * 240 * @throws CommandException thrown if failed to merge configuration 241 */ 242 protected void mergeDefaultConfig() throws CommandException { 243 Path configDefault = null; 244 try { 245 String bundleAppPathStr = conf.get(OozieClient.BUNDLE_APP_PATH); 246 Path bundleAppPath = new Path(bundleAppPathStr); 247 String user = ParamChecker.notEmpty(conf.get(OozieClient.USER_NAME), OozieClient.USER_NAME); 248 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 249 Configuration fsConf = has.createJobConf(bundleAppPath.toUri().getAuthority()); 250 FileSystem fs = has.createFileSystem(user, bundleAppPath.toUri(), fsConf); 251 252 // app path could be a directory 253 if (!fs.isFile(bundleAppPath)) { 254 configDefault = new Path(bundleAppPath, CONFIG_DEFAULT); 255 } else { 256 configDefault = new Path(bundleAppPath.getParent(), CONFIG_DEFAULT); 257 } 258 259 if (fs.exists(configDefault)) { 260 Configuration defaultConf = new XConfiguration(fs.open(configDefault)); 261 PropertiesUtils.checkDisallowedProperties(defaultConf, DISALLOWED_USER_PROPERTIES); 262 PropertiesUtils.checkDefaultDisallowedProperties(defaultConf); 263 XConfiguration.injectDefaults(defaultConf, conf); 264 } 265 else { 266 LOG.info("configDefault Doesn't exist " + configDefault); 267 } 268 PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_USER_PROPERTIES); 269 } 270 catch (IOException e) { 271 throw new CommandException(ErrorCode.E0702, e.getMessage() + " : Problem reading default config " 272 + configDefault, e); 273 } 274 catch (HadoopAccessorException e) { 275 throw new CommandException(e); 276 } 277 LOG.debug("Merged CONF :" + XmlUtils.prettyPrint(conf).toString()); 278 } 279 280 /** 281 * Read the application XML and validate against bundle Schema 282 * 283 * @return validated bundle XML 284 * @throws BundleJobException thrown if failed to read or validate xml 285 */ 286 private String readAndValidateXml() throws BundleJobException { 287 String appPath = ParamChecker.notEmpty(conf.get(OozieClient.BUNDLE_APP_PATH), OozieClient.BUNDLE_APP_PATH); 288 String bundleXml = readDefinition(appPath); 289 validateXml(bundleXml); 290 return bundleXml; 291 } 292 293 /** 294 * Read bundle definition. 295 * 296 * @param appPath application path. 297 * @param user user name. 298 * @param group group name. 299 * @return bundle definition. 300 * @throws BundleJobException thrown if the definition could not be read. 301 */ 302 protected String readDefinition(String appPath) throws BundleJobException { 303 String user = ParamChecker.notEmpty(conf.get(OozieClient.USER_NAME), OozieClient.USER_NAME); 304 //Configuration confHadoop = CoordUtils.getHadoopConf(conf); 305 try { 306 URI uri = new URI(appPath); 307 LOG.debug("user =" + user); 308 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 309 Configuration fsConf = has.createJobConf(uri.getAuthority()); 310 FileSystem fs = has.createFileSystem(user, uri, fsConf); 311 Path appDefPath = null; 312 313 // app path could be a directory 314 Path path = new Path(uri.getPath()); 315 if (!fs.isFile(path)) { 316 appDefPath = new Path(path, BUNDLE_XML_FILE); 317 } else { 318 appDefPath = path; 319 } 320 321 Reader reader = new InputStreamReader(fs.open(appDefPath)); 322 StringWriter writer = new StringWriter(); 323 IOUtils.copyCharStream(reader, writer); 324 return writer.toString(); 325 } 326 catch (IOException ex) { 327 LOG.warn("IOException :" + XmlUtils.prettyPrint(conf), ex); 328 throw new BundleJobException(ErrorCode.E1301, ex.getMessage(), ex); 329 } 330 catch (URISyntaxException ex) { 331 LOG.warn("URISyException :" + ex.getMessage()); 332 throw new BundleJobException(ErrorCode.E1302, appPath, ex.getMessage(), ex); 333 } 334 catch (HadoopAccessorException ex) { 335 throw new BundleJobException(ex); 336 } 337 catch (Exception ex) { 338 LOG.warn("Exception :", ex); 339 throw new BundleJobException(ErrorCode.E1301, ex.getMessage(), ex); 340 } 341 } 342 343 /** 344 * Validate against Bundle XSD file 345 * 346 * @param xmlContent input bundle xml 347 * @throws BundleJobException thrown if failed to validate xml 348 */ 349 private void validateXml(String xmlContent) throws BundleJobException { 350 try { 351 Validator validator = Services.get().get(SchemaService.class).getValidator(SchemaName.BUNDLE); 352 validator.validate(new StreamSource(new StringReader(xmlContent))); 353 } 354 catch (SAXException ex) { 355 LOG.warn("SAXException :", ex); 356 throw new BundleJobException(ErrorCode.E0701, ex.getMessage(), ex); 357 } 358 catch (IOException ex) { 359 LOG.warn("IOException :", ex); 360 throw new BundleJobException(ErrorCode.E0702, ex.getMessage(), ex); 361 } 362 } 363 364 /** 365 * Write a Bundle Job into database 366 * 367 * @param Bundle job bean 368 * @return job id 369 * @throws CommandException thrown if failed to store bundle job bean to db 370 */ 371 private String storeToDB(BundleJobBean bundleJob, String resolvedJobXml) throws CommandException { 372 try { 373 jobId = Services.get().get(UUIDService.class).generateId(ApplicationType.BUNDLE); 374 375 bundleJob.setId(jobId); 376 String name = XmlUtils.parseXml(bundleBean.getOrigJobXml()).getAttributeValue("name"); 377 name = ELUtils.resolveAppName(name, conf); 378 bundleJob.setAppName(name); 379 bundleJob.setAppPath(conf.get(OozieClient.BUNDLE_APP_PATH)); 380 // bundleJob.setStatus(BundleJob.Status.PREP); //This should be set in parent class. 381 bundleJob.setCreatedTime(new Date()); 382 bundleJob.setUser(conf.get(OozieClient.USER_NAME)); 383 String group = ConfigUtils.getWithDeprecatedCheck(conf, OozieClient.JOB_ACL, OozieClient.GROUP_NAME, null); 384 bundleJob.setGroup(group); 385 bundleJob.setConf(XmlUtils.prettyPrint(conf).toString()); 386 bundleJob.setJobXml(resolvedJobXml); 387 Element jobElement = XmlUtils.parseXml(resolvedJobXml); 388 Element controlsElement = jobElement.getChild("controls", jobElement.getNamespace()); 389 if (controlsElement != null) { 390 Element kickoffTimeElement = controlsElement.getChild("kick-off-time", jobElement.getNamespace()); 391 if (kickoffTimeElement != null && !kickoffTimeElement.getValue().isEmpty()) { 392 Date kickoffTime = DateUtils.parseDateOozieTZ(kickoffTimeElement.getValue()); 393 bundleJob.setKickoffTime(kickoffTime); 394 } 395 } 396 bundleJob.setLastModifiedTime(new Date()); 397 398 if (!dryrun) { 399 BundleJobQueryExecutor.getInstance().insert(bundleJob); 400 } 401 } 402 catch (Exception ex) { 403 throw new CommandException(ErrorCode.E1301, ex.getMessage(), ex); 404 } 405 return jobId; 406 } 407 408 /* (non-Javadoc) 409 * @see org.apache.oozie.command.TransitionXCommand#getJob() 410 */ 411 @Override 412 public Job getJob() { 413 return bundleBean; 414 } 415 416 public static ELEvaluator createELEvaluatorForGroup(Configuration conf, String group) { 417 ELEvaluator eval = Services.get().get(ELService.class).createEvaluator(group); 418 setConfigToEval(eval, conf); 419 return eval; 420 } 421 422 private static void setConfigToEval(ELEvaluator eval, Configuration conf) { 423 for (Map.Entry<String, String> entry : conf) { 424 eval.setVariable(entry.getKey(), entry.getValue().trim()); 425 } 426 } 427 428 /** 429 * Resolve job xml with conf 430 * 431 * @param bundleXml bundle job xml 432 * @param conf job configuration 433 * @return resolved job xml 434 * @throws BundleJobException thrown if failed to resolve variables 435 */ 436 private String resolvedVarsandFunctions(String bundleXml, Configuration conf) throws BundleJobException { 437 ELEvaluator eval; 438 try { 439 eval = createELEvaluatorForGroup(conf, "bundle-submit"); 440 return eval.evaluate(bundleXml, String.class); 441 } 442 catch (Exception e) { 443 throw new BundleJobException(ErrorCode.E1004, e.getMessage(), e); 444 } 445 } 446 447 /** 448 * Create ELEvaluator 449 * 450 * @param conf job configuration 451 * @return ELEvaluator the evaluator for el function 452 * @throws BundleJobException thrown if failed to create evaluator 453 */ 454 public ELEvaluator createEvaluator(Configuration conf) throws BundleJobException { 455 ELEvaluator eval; 456 ELEvaluator.Context context; 457 try { 458 context = new ELEvaluator.Context(); 459 eval = new ELEvaluator(context); 460 for (Map.Entry<String, String> entry : conf) { 461 eval.setVariable(entry.getKey(), entry.getValue()); 462 } 463 } 464 catch (Exception e) { 465 throw new BundleJobException(ErrorCode.E1004, e.getMessage(), e); 466 } 467 return eval; 468 } 469 470 /** 471 * Verify the uniqueness of coordinator names 472 * 473 * @param resolved job xml 474 * @throws CommandException thrown if failed to verify the uniqueness of coordinator names 475 */ 476 @SuppressWarnings("unchecked") 477 private Void verifyCoordNameUnique(String resolvedJobXml) throws CommandException { 478 Set<String> set = new HashSet<String>(); 479 try { 480 Element bAppXml = XmlUtils.parseXml(resolvedJobXml); 481 List<Element> coordElems = bAppXml.getChildren("coordinator", bAppXml.getNamespace()); 482 for (Element elem : coordElems) { 483 Attribute name = elem.getAttribute("name"); 484 if (name != null) { 485 String coordName = name.getValue(); 486 try { 487 coordName = ELUtils.resolveAppName(name.getValue(), conf); 488 } 489 catch (Exception e) { 490 throw new CommandException(ErrorCode.E1321, e.getMessage(), e); 491 } 492 if (set.contains(coordName)) { 493 throw new CommandException(ErrorCode.E1304, name); 494 } 495 set.add(coordName); 496 } 497 else { 498 throw new CommandException(ErrorCode.E1305); 499 } 500 } 501 } 502 catch (JDOMException jex) { 503 throw new CommandException(ErrorCode.E1301, jex.getMessage(), jex); 504 } 505 506 return null; 507 } 508 509 /* (non-Javadoc) 510 * @see org.apache.oozie.command.TransitionXCommand#updateJob() 511 */ 512 @Override 513 public void updateJob() throws CommandException { 514 } 515 516 @Override 517 public void performWrites() throws CommandException { 518 } 519}