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.wf; 020 021import java.io.IOException; 022import java.net.URI; 023import java.net.URISyntaxException; 024import java.util.ArrayList; 025import java.util.Collection; 026import java.util.Date; 027import java.util.HashMap; 028import java.util.HashSet; 029import java.util.List; 030import java.util.Map; 031import java.util.Set; 032 033import org.apache.hadoop.conf.Configuration; 034import org.apache.hadoop.fs.FileSystem; 035import org.apache.hadoop.fs.Path; 036import org.apache.oozie.AppType; 037import org.apache.oozie.ErrorCode; 038import org.apache.oozie.WorkflowActionBean; 039import org.apache.oozie.WorkflowJobBean; 040import org.apache.oozie.action.oozie.SubWorkflowActionExecutor; 041import org.apache.oozie.client.OozieClient; 042import org.apache.oozie.client.WorkflowAction; 043import org.apache.oozie.client.WorkflowJob; 044import org.apache.oozie.client.rest.JsonBean; 045import org.apache.oozie.command.CommandException; 046import org.apache.oozie.command.PreconditionException; 047import org.apache.oozie.executor.jpa.JPAExecutorException; 048import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor; 049import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor.WorkflowActionQuery; 050import org.apache.oozie.executor.jpa.BatchQueryExecutor; 051import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor; 052import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry; 053import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor.WorkflowJobQuery; 054import org.apache.oozie.service.ConfigurationService; 055import org.apache.oozie.service.DagXLogInfoService; 056import org.apache.oozie.service.HadoopAccessorException; 057import org.apache.oozie.service.HadoopAccessorService; 058import org.apache.oozie.service.Services; 059import org.apache.oozie.service.UUIDService; 060import org.apache.oozie.service.WorkflowAppService; 061import org.apache.oozie.service.WorkflowStoreService; 062import org.apache.oozie.sla.SLAOperations; 063import org.apache.oozie.sla.service.SLAService; 064import org.apache.oozie.util.ConfigUtils; 065import org.apache.oozie.util.ELEvaluator; 066import org.apache.oozie.util.ELUtils; 067import org.apache.oozie.util.InstrumentUtils; 068import org.apache.oozie.util.LogUtils; 069import org.apache.oozie.util.ParamChecker; 070import org.apache.oozie.util.PropertiesUtils; 071import org.apache.oozie.util.XConfiguration; 072import org.apache.oozie.util.XLog; 073import org.apache.oozie.util.XmlUtils; 074import org.apache.oozie.workflow.WorkflowApp; 075import org.apache.oozie.workflow.WorkflowException; 076import org.apache.oozie.workflow.WorkflowInstance; 077import org.apache.oozie.workflow.WorkflowLib; 078import org.apache.oozie.workflow.lite.NodeHandler; 079import org.jdom.Element; 080import org.jdom.JDOMException; 081 082/** 083 * This is a RerunXCommand which is used for rerunn. 084 * 085 */ 086public class ReRunXCommand extends WorkflowXCommand<Void> { 087 private final String jobId; 088 private Configuration conf; 089 private final Set<String> nodesToSkip = new HashSet<String>(); 090 public static final String TO_SKIP = "TO_SKIP"; 091 private WorkflowJobBean wfBean; 092 private List<WorkflowActionBean> actions; 093 private List<UpdateEntry> updateList = new ArrayList<UpdateEntry>(); 094 private List<JsonBean> deleteList = new ArrayList<JsonBean>(); 095 096 private static final Set<String> DISALLOWED_USER_PROPERTIES = new HashSet<String>(); 097 public static final String DISABLE_CHILD_RERUN = "oozie.wf.rerun.disablechild"; 098 099 static { 100 String[] badUserProps = { PropertiesUtils.DAYS, PropertiesUtils.HOURS, PropertiesUtils.MINUTES, 101 PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, PropertiesUtils.TB, PropertiesUtils.PB, 102 PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN, 103 PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS }; 104 PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_USER_PROPERTIES); 105 } 106 107 public ReRunXCommand(String jobId, Configuration conf) { 108 super("rerun", "rerun", 1); 109 this.jobId = ParamChecker.notEmpty(jobId, "jobId"); 110 this.conf = ParamChecker.notNull(conf, "conf"); 111 } 112 113 @Override 114 protected void setLogInfo() { 115 LogUtils.setLogInfo(jobId); 116 } 117 118 /* (non-Javadoc) 119 * @see org.apache.oozie.command.XCommand#execute() 120 */ 121 @Override 122 protected Void execute() throws CommandException { 123 setupReRun(); 124 startWorkflow(jobId); 125 return null; 126 } 127 128 private void startWorkflow(String jobId) throws CommandException { 129 new StartXCommand(jobId).call(); 130 } 131 132 private void setupReRun() throws CommandException { 133 InstrumentUtils.incrJobCounter(getName(), 1, getInstrumentation()); 134 LogUtils.setLogInfo(wfBean); 135 WorkflowInstance oldWfInstance = this.wfBean.getWorkflowInstance(); 136 WorkflowInstance newWfInstance; 137 String appPath = null; 138 139 WorkflowAppService wps = Services.get().get(WorkflowAppService.class); 140 try { 141 XLog.Info.get().setParameter(DagXLogInfoService.TOKEN, conf.get(OozieClient.LOG_TOKEN)); 142 WorkflowApp app = wps.parseDef(conf, null); 143 XConfiguration protoActionConf = wps.createProtoActionConf(conf, true); 144 WorkflowLib workflowLib = Services.get().get(WorkflowStoreService.class).getWorkflowLibWithNoDB(); 145 146 appPath = conf.get(OozieClient.APP_PATH); 147 URI uri = new URI(appPath); 148 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 149 Configuration fsConf = has.createJobConf(uri.getAuthority()); 150 FileSystem fs = has.createFileSystem(wfBean.getUser(), uri, fsConf); 151 152 Path configDefault = null; 153 // app path could be a directory 154 Path path = new Path(uri.getPath()); 155 if (!fs.isFile(path)) { 156 configDefault = new Path(path, SubmitXCommand.CONFIG_DEFAULT); 157 } 158 else { 159 configDefault = new Path(path.getParent(), SubmitXCommand.CONFIG_DEFAULT); 160 } 161 162 if (fs.exists(configDefault)) { 163 Configuration defaultConf = new XConfiguration(fs.open(configDefault)); 164 PropertiesUtils.checkDisallowedProperties(defaultConf, DISALLOWED_USER_PROPERTIES); 165 PropertiesUtils.checkDefaultDisallowedProperties(defaultConf); 166 XConfiguration.injectDefaults(defaultConf, conf); 167 } 168 169 PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_USER_PROPERTIES); 170 171 // Resolving all variables in the job properties. This ensures the Hadoop Configuration semantics are 172 // preserved. The Configuration.get function within XConfiguration.resolve() works recursively to get the 173 // final value corresponding to a key in the map Resetting the conf to contain all the resolved values is 174 // necessary to ensure propagation of Oozie properties to Hadoop calls downstream 175 conf = ((XConfiguration) conf).resolve(); 176 177 try { 178 newWfInstance = workflowLib.createInstance(app, conf, jobId); 179 } 180 catch (WorkflowException e) { 181 throw new CommandException(e); 182 } 183 String appName = ELUtils.resolveAppName(app.getName(), conf); 184 if (SLAService.isEnabled()) { 185 Element wfElem = XmlUtils.parseXml(app.getDefinition()); 186 ELEvaluator evalSla = SubmitXCommand.createELEvaluatorForGroup(conf, "wf-sla-submit"); 187 Element eSla = XmlUtils.getSLAElement(wfElem); 188 String jobSlaXml = null; 189 if (eSla != null) { 190 jobSlaXml = SubmitXCommand.resolveSla(eSla, evalSla); 191 } 192 writeSLARegistration(wfElem, jobSlaXml, newWfInstance.getId(), 193 conf.get(SubWorkflowActionExecutor.PARENT_ID), conf.get(OozieClient.USER_NAME), appName, 194 evalSla); 195 } 196 wfBean.setAppName(appName); 197 wfBean.setProtoActionConf(protoActionConf.toXmlString()); 198 } 199 catch (WorkflowException ex) { 200 throw new CommandException(ex); 201 } 202 catch (IOException ex) { 203 throw new CommandException(ErrorCode.E0803, ex.getMessage(), ex); 204 } 205 catch (HadoopAccessorException ex) { 206 throw new CommandException(ex); 207 } 208 catch (URISyntaxException ex) { 209 throw new CommandException(ErrorCode.E0711, appPath, ex.getMessage(), ex); 210 } 211 catch (Exception ex) { 212 throw new CommandException(ErrorCode.E1007, ex.getMessage(), ex); 213 } 214 215 for (int i = 0; i < actions.size(); i++) { 216 // Skipping to delete the sub workflow when rerun failed node option has been provided. As same 217 // action will be used to rerun the job. 218 if (!nodesToSkip.contains(actions.get(i).getName()) && 219 !(conf.getBoolean(OozieClient.RERUN_FAIL_NODES, false) && 220 SubWorkflowActionExecutor.ACTION_TYPE.equals(actions.get(i).getType()))) { 221 deleteList.add(actions.get(i)); 222 LOG.info("Deleting Action[{0}] for re-run", actions.get(i).getId()); 223 } 224 else { 225 copyActionData(newWfInstance, oldWfInstance); 226 } 227 } 228 229 wfBean.setAppPath(conf.get(OozieClient.APP_PATH)); 230 wfBean.setConf(XmlUtils.prettyPrint(conf).toString()); 231 wfBean.setLogToken(conf.get(OozieClient.LOG_TOKEN, "")); 232 wfBean.setUser(conf.get(OozieClient.USER_NAME)); 233 String group = ConfigUtils.getWithDeprecatedCheck(conf, OozieClient.JOB_ACL, OozieClient.GROUP_NAME, null); 234 wfBean.setGroup(group); 235 wfBean.setExternalId(conf.get(OozieClient.EXTERNAL_ID)); 236 wfBean.setEndTime(null); 237 wfBean.setRun(wfBean.getRun() + 1); 238 wfBean.setStatus(WorkflowJob.Status.PREP); 239 wfBean.setWorkflowInstance(newWfInstance); 240 241 try { 242 wfBean.setLastModifiedTime(new Date()); 243 updateList.add(new UpdateEntry<WorkflowJobQuery>(WorkflowJobQuery.UPDATE_WORKFLOW_RERUN, wfBean)); 244 // call JPAExecutor to do the bulk writes 245 BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(null, updateList, deleteList); 246 } 247 catch (JPAExecutorException je) { 248 throw new CommandException(je); 249 } 250 finally { 251 updateParentIfNecessary(wfBean); 252 } 253 254 } 255 256 @SuppressWarnings("unchecked") 257 private void writeSLARegistration(Element wfElem, String jobSlaXml, String id, String parentId, String user, 258 String appName, ELEvaluator evalSla) throws JDOMException, CommandException { 259 if (jobSlaXml != null && jobSlaXml.length() > 0) { 260 Element eSla = XmlUtils.parseXml(jobSlaXml); 261 // insert into new table 262 SLAOperations.createSlaRegistrationEvent(eSla, jobId, parentId, AppType.WORKFLOW_JOB, user, appName, LOG, 263 true); 264 } 265 // Add sla for wf actions 266 for (Element action : (List<Element>) wfElem.getChildren("action", wfElem.getNamespace())) { 267 Element actionSla = XmlUtils.getSLAElement(action); 268 if (actionSla != null) { 269 String actionSlaXml = SubmitXCommand.resolveSla(actionSla, evalSla); 270 actionSla = XmlUtils.parseXml(actionSlaXml); 271 if (!nodesToSkip.contains(action.getAttributeValue("name"))) { 272 String actionId = Services.get().get(UUIDService.class) 273 .generateChildId(jobId, action.getAttributeValue("name") + ""); 274 SLAOperations.createSlaRegistrationEvent(actionSla, actionId, jobId, AppType.WORKFLOW_ACTION, user, 275 appName, LOG, true); 276 } 277 } 278 } 279 280 } 281 282 /** 283 * Loading the Wfjob and workflow actions. Parses the config and adds the nodes that are to be skipped to the 284 * skipped node list 285 * 286 * @throws CommandException 287 */ 288 @Override 289 protected void eagerLoadState() throws CommandException { 290 try { 291 this.wfBean = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW_STATUS, this.jobId); 292 this.actions = WorkflowActionQueryExecutor.getInstance().getList( 293 WorkflowActionQuery.GET_ACTIONS_FOR_WORKFLOW_RERUN, this.jobId); 294 295 if (conf != null) { 296 if (conf.getBoolean(OozieClient.RERUN_FAIL_NODES, false) == false) { //Rerun with skipNodes 297 Collection<String> skipNodes = conf.getStringCollection(OozieClient.RERUN_SKIP_NODES); 298 for (String str : skipNodes) { 299 // trimming is required 300 nodesToSkip.add(str.trim()); 301 } 302 LOG.debug("Skipnode size :" + nodesToSkip.size()); 303 } 304 else { 305 for (WorkflowActionBean action : actions) { // Rerun from failed nodes 306 if (action.getStatus() == WorkflowAction.Status.OK) { 307 nodesToSkip.add(action.getName()); 308 } 309 } 310 LOG.debug("Skipnode size are to rerun from FAIL nodes :" + nodesToSkip.size()); 311 } 312 StringBuilder tmp = new StringBuilder(); 313 for (String node : nodesToSkip) { 314 tmp.append(node).append(","); 315 } 316 LOG.debug("SkipNode List :" + tmp); 317 } 318 } 319 catch (Exception ex) { 320 throw new CommandException(ErrorCode.E0603, ex.getMessage(), ex); 321 } 322 } 323 324 /** 325 * Checks the pre-conditions that are required for workflow to recover - Last run of Workflow should be completed - 326 * The nodes that are to be skipped are to be completed successfully in the base run. 327 * 328 * @throws org.apache.oozie.command.CommandException,PreconditionException On failure of pre-conditions 329 */ 330 @Override 331 protected void eagerVerifyPrecondition() throws CommandException, PreconditionException { 332 // Throwing error if parent exist and same workflow trying to rerun, when running child workflow disabled 333 // through conf. 334 if (wfBean.getParentId() != null && !conf.getBoolean(SubWorkflowActionExecutor.SUBWORKFLOW_RERUN, false) 335 && ConfigurationService.getBoolean(DISABLE_CHILD_RERUN)) { 336 throw new PreconditionException(ErrorCode.E0755, " Rerun is not allowed through child workflow, please" + 337 " re-run through the parent " + wfBean.getParentId()); 338 } 339 340 if (!(wfBean.getStatus().equals(WorkflowJob.Status.FAILED) 341 || wfBean.getStatus().equals(WorkflowJob.Status.KILLED) || wfBean.getStatus().equals( 342 WorkflowJob.Status.SUCCEEDED))) { 343 throw new CommandException(ErrorCode.E0805, wfBean.getStatus()); 344 } 345 Set<String> unmachedNodes = new HashSet<String>(nodesToSkip); 346 for (WorkflowActionBean action : actions) { 347 if (nodesToSkip.contains(action.getName())) { 348 if (!action.getStatus().equals(WorkflowAction.Status.OK) 349 && !action.getStatus().equals(WorkflowAction.Status.ERROR)) { 350 throw new CommandException(ErrorCode.E0806, action.getName()); 351 } 352 unmachedNodes.remove(action.getName()); 353 } 354 } 355 if (unmachedNodes.size() > 0) { 356 StringBuilder sb = new StringBuilder(); 357 String separator = ""; 358 for (String s : unmachedNodes) { 359 sb.append(separator).append(s); 360 separator = ","; 361 } 362 throw new CommandException(ErrorCode.E0807, sb); 363 } 364 } 365 366 /** 367 * Copys the variables for skipped nodes from the old wfInstance to new one. 368 * 369 * @param newWfInstance : Source WF instance object 370 * @param oldWfInstance : Update WF instance 371 */ 372 private void copyActionData(WorkflowInstance newWfInstance, WorkflowInstance oldWfInstance) { 373 Map<String, String> oldVars = new HashMap<String, String>(); 374 Map<String, String> newVars = new HashMap<String, String>(); 375 oldVars = oldWfInstance.getAllVars(); 376 for (String var : oldVars.keySet()) { 377 String actionName = var.split(WorkflowInstance.NODE_VAR_SEPARATOR)[0]; 378 if (nodesToSkip.contains(actionName)) { 379 newVars.put(var, oldVars.get(var)); 380 } 381 } 382 for (String node : nodesToSkip) { 383 // Setting the TO_SKIP variable to true. This will be used by 384 // SignalCommand and LiteNodeHandler to skip the action. 385 newVars.put(node + WorkflowInstance.NODE_VAR_SEPARATOR + TO_SKIP, "true"); 386 String visitedFlag = NodeHandler.getLoopFlag(node); 387 // Removing the visited flag so that the action won't be considered 388 // a loop. 389 if (newVars.containsKey(visitedFlag)) { 390 newVars.remove(visitedFlag); 391 } 392 } 393 newWfInstance.setAllVars(newVars); 394 } 395 396 /* (non-Javadoc) 397 * @see org.apache.oozie.command.XCommand#getEntityKey() 398 */ 399 @Override 400 public String getEntityKey() { 401 return this.jobId; 402 } 403 404 /* (non-Javadoc) 405 * @see org.apache.oozie.command.XCommand#isLockRequired() 406 */ 407 @Override 408 protected boolean isLockRequired() { 409 return true; 410 } 411 412 /* (non-Javadoc) 413 * @see org.apache.oozie.command.XCommand#loadState() 414 */ 415 @Override 416 protected void loadState() throws CommandException { 417 try { 418 this.wfBean = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW_RERUN, this.jobId); 419 this.actions = WorkflowActionQueryExecutor.getInstance().getList( 420 WorkflowActionQuery.GET_ACTIONS_FOR_WORKFLOW_RERUN, this.jobId); 421 } 422 catch (JPAExecutorException jpe) { 423 throw new CommandException(jpe); 424 } 425 } 426 427 /* (non-Javadoc) 428 * @see org.apache.oozie.command.XCommand#verifyPrecondition() 429 */ 430 @Override 431 protected void verifyPrecondition() throws CommandException, PreconditionException { 432 eagerVerifyPrecondition(); 433 } 434}