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.StringReader; 023import java.util.Date; 024import java.util.HashMap; 025import java.util.List; 026import java.util.Map; 027import java.util.Map.Entry; 028 029import org.apache.hadoop.conf.Configuration; 030import org.apache.oozie.BundleActionBean; 031import org.apache.oozie.BundleJobBean; 032import org.apache.oozie.ErrorCode; 033import org.apache.oozie.XException; 034import org.apache.oozie.action.hadoop.OozieJobInfo; 035import org.apache.oozie.client.Job; 036import org.apache.oozie.client.OozieClient; 037import org.apache.oozie.client.rest.JsonBean; 038import org.apache.oozie.command.CommandException; 039import org.apache.oozie.command.PreconditionException; 040import org.apache.oozie.command.StartTransitionXCommand; 041import org.apache.oozie.command.coord.CoordSubmitXCommand; 042import org.apache.oozie.executor.jpa.BatchQueryExecutor; 043import org.apache.oozie.executor.jpa.BundleJobQueryExecutor; 044import org.apache.oozie.executor.jpa.BundleJobQueryExecutor.BundleJobQuery; 045import org.apache.oozie.executor.jpa.JPAExecutorException; 046import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry; 047import org.apache.oozie.util.ConfigUtils; 048import org.apache.oozie.util.ELUtils; 049import org.apache.oozie.util.JobUtils; 050import org.apache.oozie.util.LogUtils; 051import org.apache.oozie.util.ParamChecker; 052import org.apache.oozie.util.XConfiguration; 053import org.apache.oozie.util.XmlUtils; 054import org.jdom.Attribute; 055import org.jdom.Element; 056import org.jdom.JDOMException; 057 058/** 059 * The command to start Bundle job 060 */ 061public class BundleStartXCommand extends StartTransitionXCommand { 062 private final String jobId; 063 private BundleJobBean bundleJob; 064 065 /** 066 * The constructor for class {@link BundleStartXCommand} 067 * 068 * @param jobId the bundle job id 069 */ 070 public BundleStartXCommand(String jobId) { 071 super("bundle_start", "bundle_start", 1); 072 this.jobId = ParamChecker.notEmpty(jobId, "jobId"); 073 } 074 075 /** 076 * The constructor for class {@link BundleStartXCommand} 077 * 078 * @param jobId the bundle job id 079 * @param dryrun true if dryrun is enable 080 */ 081 public BundleStartXCommand(String jobId, boolean dryrun) { 082 super("bundle_start", "bundle_start", 1, dryrun); 083 this.jobId = ParamChecker.notEmpty(jobId, "jobId"); 084 } 085 086 /* (non-Javadoc) 087 * @see org.apache.oozie.command.XCommand#getEntityKey() 088 */ 089 @Override 090 public String getEntityKey() { 091 return jobId; 092 } 093 094 @Override 095 public String getKey() { 096 return getName() + "_" + jobId; 097 } 098 099 /* (non-Javadoc) 100 * @see org.apache.oozie.command.XCommand#isLockRequired() 101 */ 102 @Override 103 protected boolean isLockRequired() { 104 return true; 105 } 106 107 @Override 108 protected void verifyPrecondition() throws CommandException, PreconditionException { 109 if (bundleJob.getStatus() != Job.Status.PREP) { 110 String msg = "Bundle " + bundleJob.getId() + " is not in PREP status. It is in : " + bundleJob.getStatus(); 111 LOG.info(msg); 112 throw new PreconditionException(ErrorCode.E1100, msg); 113 } 114 } 115 116 @Override 117 public void loadState() throws CommandException { 118 try { 119 this.bundleJob = BundleJobQueryExecutor.getInstance().get(BundleJobQuery.GET_BUNDLE_JOB, jobId); 120 LogUtils.setLogInfo(bundleJob); 121 super.setJob(bundleJob); 122 } 123 catch (XException ex) { 124 throw new CommandException(ex); 125 } 126 } 127 128 /* (non-Javadoc) 129 * @see org.apache.oozie.command.StartTransitionXCommand#StartChildren() 130 */ 131 @Override 132 public void StartChildren() throws CommandException { 133 LOG.debug("Started coord jobs for the bundle=[{0}]", jobId); 134 insertBundleActions(); 135 startCoordJobs(); 136 LOG.debug("Ended coord jobs for the bundle=[{0}]", jobId); 137 } 138 139 /* (non-Javadoc) 140 * @see org.apache.oozie.command.TransitionXCommand#notifyParent() 141 */ 142 @Override 143 public void notifyParent() { 144 } 145 146 /* (non-Javadoc) 147 * @see org.apache.oozie.command.StartTransitionXCommand#performWrites() 148 */ 149 @Override 150 public void performWrites() throws CommandException { 151 try { 152 BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, updateList, null); 153 } 154 catch (JPAExecutorException e) { 155 throw new CommandException(e); 156 } 157 } 158 159 /** 160 * Insert bundle actions 161 * 162 * @throws CommandException thrown if failed to create bundle actions 163 */ 164 @SuppressWarnings("unchecked") 165 private void insertBundleActions() throws CommandException { 166 if (bundleJob != null) { 167 Map<String, Boolean> map = new HashMap<String, Boolean>(); 168 try { 169 Element bAppXml = XmlUtils.parseXml(bundleJob.getJobXml()); 170 List<Element> coordElems = bAppXml.getChildren("coordinator", bAppXml.getNamespace()); 171 for (Element elem : coordElems) { 172 Attribute name = elem.getAttribute("name"); 173 Attribute critical = elem.getAttribute("critical"); 174 if (name != null) { 175 if (map.containsKey(name.getValue())) { 176 throw new CommandException(ErrorCode.E1304, name); 177 } 178 boolean isCritical = false; 179 if (critical != null && Boolean.parseBoolean(critical.getValue())) { 180 isCritical = true; 181 } 182 map.put(name.getValue(), isCritical); 183 } 184 else { 185 throw new CommandException(ErrorCode.E1305); 186 } 187 } 188 } 189 catch (JDOMException jex) { 190 throw new CommandException(ErrorCode.E1301, jex.getMessage(), jex); 191 } 192 193 // if there is no coordinator for this bundle, failed it. 194 if (map.isEmpty()) { 195 bundleJob.setStatus(Job.Status.FAILED); 196 bundleJob.resetPending(); 197 try { 198 BundleJobQueryExecutor.getInstance().executeUpdate(BundleJobQuery.UPDATE_BUNDLE_JOB_STATUS_PENDING, bundleJob); 199 } 200 catch (JPAExecutorException jex) { 201 throw new CommandException(jex); 202 } 203 204 LOG.debug("No coord jobs for the bundle=[{0}], failed it!!", jobId); 205 throw new CommandException(ErrorCode.E1318, jobId); 206 } 207 208 for (Entry<String, Boolean> coordName : map.entrySet()) { 209 BundleActionBean action = createBundleAction(jobId, coordName.getKey(), coordName.getValue()); 210 insertList.add(action); 211 } 212 } 213 else { 214 throw new CommandException(ErrorCode.E0604, jobId); 215 } 216 } 217 218 private BundleActionBean createBundleAction(String jobId, String coordName, boolean isCritical) { 219 BundleActionBean action = new BundleActionBean(); 220 action.setBundleActionId(jobId + "_" + coordName); 221 action.setBundleId(jobId); 222 action.setCoordName(coordName); 223 action.setStatus(Job.Status.PREP); 224 action.setLastModifiedTime(new Date()); 225 if (isCritical) { 226 action.setCritical(); 227 } 228 else { 229 action.resetCritical(); 230 } 231 return action; 232 } 233 234 /** 235 * Start Coord Jobs 236 * 237 * @throws CommandException thrown if failed to start coord jobs 238 */ 239 @SuppressWarnings("unchecked") 240 private void startCoordJobs() throws CommandException { 241 if (bundleJob != null) { 242 try { 243 Element bAppXml = XmlUtils.parseXml(bundleJob.getJobXml()); 244 List<Element> coordElems = bAppXml.getChildren("coordinator", bAppXml.getNamespace()); 245 for (Element coordElem : coordElems) { 246 Attribute name = coordElem.getAttribute("name"); 247 248 Configuration coordConf = mergeConfig(coordElem); 249 coordConf.set(OozieClient.BUNDLE_ID, jobId); 250 if (OozieJobInfo.isJobInfoEnabled()) { 251 coordConf.set(OozieJobInfo.BUNDLE_NAME, bundleJob.getAppName()); 252 } 253 String coordName=name.getValue(); 254 try { 255 coordName = ELUtils.resolveAppName(coordName, coordConf); 256 } 257 catch (Exception e) { 258 throw new CommandException(ErrorCode.E1321, e.getMessage(), e); 259 260 } 261 queue(new CoordSubmitXCommand(coordConf, bundleJob.getId(), name.getValue())); 262 263 } 264 updateBundleAction(); 265 } 266 catch (JDOMException jex) { 267 throw new CommandException(ErrorCode.E1301, jex.getMessage(), jex); 268 } 269 catch (JPAExecutorException je) { 270 throw new CommandException(je); 271 } 272 } 273 else { 274 throw new CommandException(ErrorCode.E0604, jobId); 275 } 276 } 277 278 private void updateBundleAction() throws JPAExecutorException { 279 for(JsonBean bAction : insertList) { 280 BundleActionBean action = (BundleActionBean) bAction; 281 action.incrementAndGetPending(); 282 action.setLastModifiedTime(new Date()); 283 } 284 } 285 286 /** 287 * Merge Bundle job config and the configuration from the coord job to pass 288 * to Coord Engine 289 * 290 * @param coordElem the coordinator configuration 291 * @return Configuration merged configuration 292 * @throws CommandException thrown if failed to merge configuration 293 */ 294 private Configuration mergeConfig(Element coordElem) throws CommandException { 295 String jobConf = bundleJob.getConf(); 296 // Step 1: runConf = jobConf 297 Configuration runConf = null; 298 try { 299 runConf = new XConfiguration(new StringReader(jobConf)); 300 } 301 catch (IOException e1) { 302 LOG.warn("Configuration parse error in:" + jobConf); 303 throw new CommandException(ErrorCode.E1306, e1.getMessage(), e1); 304 } 305 // Step 2: Merge local properties into runConf 306 // extract 'property' tags under 'configuration' block in the coordElem 307 // convert Element to XConfiguration 308 Element localConfigElement = coordElem.getChild("configuration", coordElem.getNamespace()); 309 310 if (localConfigElement != null) { 311 String strConfig = XmlUtils.prettyPrint(localConfigElement).toString(); 312 Configuration localConf; 313 try { 314 localConf = new XConfiguration(new StringReader(strConfig)); 315 } 316 catch (IOException e1) { 317 LOG.warn("Configuration parse error in:" + strConfig); 318 throw new CommandException(ErrorCode.E1307, e1.getMessage(), e1); 319 } 320 321 // copy configuration properties in the coordElem to the runConf 322 XConfiguration.copy(localConf, runConf); 323 324 ConfigUtils.checkAndSetDisallowedProperties(runConf, 325 bundleJob.getUser(), 326 new CommandException(ErrorCode.E1303, 327 String.format("%s=%s", OozieClient.USER_NAME, runConf.get(OozieClient.USER_NAME)), 328 bundleJob.getUser()), 329 true); 330 } 331 332 // Step 3: Extract value of 'app-path' in coordElem, save it as a 333 // new property called 'oozie.coord.application.path', and normalize. 334 String appPath = coordElem.getChild("app-path", coordElem.getNamespace()).getValue(); 335 runConf.set(OozieClient.COORDINATOR_APP_PATH, appPath); 336 // Normalize coordinator appPath here; 337 try { 338 JobUtils.normalizeAppPath(runConf.get(OozieClient.USER_NAME), runConf.get(OozieClient.GROUP_NAME), runConf); 339 } 340 catch (IOException e) { 341 throw new CommandException(ErrorCode.E1001, runConf.get(OozieClient.COORDINATOR_APP_PATH)); 342 } 343 return runConf; 344 } 345 346 /* (non-Javadoc) 347 * @see org.apache.oozie.command.TransitionXCommand#getJob() 348 */ 349 @Override 350 public Job getJob() { 351 return bundleJob; 352 } 353 354 /* (non-Javadoc) 355 * @see org.apache.oozie.command.TransitionXCommand#updateJob() 356 */ 357 @Override 358 public void updateJob() throws CommandException { 359 updateList.add(new UpdateEntry<BundleJobQuery>(BundleJobQuery.UPDATE_BUNDLE_JOB_STATUS_PENDING, bundleJob)); 360 } 361}