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.io.StringReader; 023import java.net.URI; 024import java.net.URISyntaxException; 025import java.util.Date; 026import java.util.HashMap; 027import java.util.Map; 028import java.util.Properties; 029import java.util.Set; 030 031import org.apache.hadoop.conf.Configuration; 032import org.apache.hadoop.fs.FileSystem; 033import org.apache.hadoop.fs.Path; 034import org.apache.oozie.DagELFunctions; 035import org.apache.oozie.ErrorCode; 036import org.apache.oozie.WorkflowActionBean; 037import org.apache.oozie.WorkflowJobBean; 038import org.apache.oozie.action.ActionExecutor; 039import org.apache.oozie.client.Job; 040import org.apache.oozie.client.WorkflowAction; 041import org.apache.oozie.client.WorkflowJob; 042import org.apache.oozie.command.CommandException; 043import org.apache.oozie.service.CallbackService; 044import org.apache.oozie.service.ELService; 045import org.apache.oozie.service.HadoopAccessorException; 046import org.apache.oozie.service.HadoopAccessorService; 047import org.apache.oozie.service.JPAService; 048import org.apache.oozie.service.LiteWorkflowStoreService; 049import org.apache.oozie.service.Services; 050import org.apache.oozie.util.ELEvaluator; 051import org.apache.oozie.util.InstrumentUtils; 052import org.apache.oozie.util.Instrumentation; 053import org.apache.oozie.util.XConfiguration; 054import org.apache.oozie.workflow.WorkflowException; 055import org.apache.oozie.workflow.WorkflowInstance; 056import org.apache.oozie.workflow.lite.LiteWorkflowInstance; 057 058/** 059 * Base class for Action execution commands. Provides common functionality to handle different types of errors while 060 * attempting to start or end an action. 061 */ 062public abstract class ActionXCommand<T> extends WorkflowXCommand<T> { 063 private static final String INSTRUMENTATION_GROUP = "action.executors"; 064 065 protected static final String RECOVERY_ID_SEPARATOR = "@"; 066 067 public ActionXCommand(String name, String type, int priority) { 068 super(name, type, priority); 069 } 070 071 /** 072 * Takes care of Transient failures. Sets the action status to retry and increments the retry count if not enough 073 * attempts have been made. Otherwise returns false. 074 * 075 * @param context the execution context. 076 * @param executor the executor instance being used. 077 * @param status the status to be set for the action. 078 * @return true if the action is scheduled for another retry. false if the number of retries has exceeded the 079 * maximum number of configured retries. 080 * @throws CommandException thrown if unable to handle transient 081 */ 082 protected boolean handleTransient(ActionExecutor.Context context, ActionExecutor executor, 083 WorkflowAction.Status status) throws CommandException { 084 LOG.debug("Attempting to retry"); 085 ActionExecutorContext aContext = (ActionExecutorContext) context; 086 WorkflowActionBean action = (WorkflowActionBean) aContext.getAction(); 087 incrActionErrorCounter(action.getType(), "transient", 1); 088 089 int actionRetryCount = action.getRetries(); 090 if (actionRetryCount >= executor.getMaxRetries()) { 091 LOG.warn("Exceeded max retry count [{0}]. Suspending Job", executor.getMaxRetries()); 092 return false; 093 } 094 else { 095 action.setStatus(status); 096 action.setPending(); 097 action.incRetries(); 098 long retryDelayMillis = getRetryDelay(actionRetryCount, executor.getRetryInterval(), executor.getRetryPolicy()); 099 action.setPendingAge(new Date(System.currentTimeMillis() + retryDelayMillis)); 100 LOG.info("Next Retry, Attempt Number [{0}] in [{1}] milliseconds", actionRetryCount + 1, retryDelayMillis); 101 this.resetUsed(); 102 queueCommandForTransientFailure(retryDelayMillis); 103 return true; 104 } 105 } 106 107 protected void queueCommandForTransientFailure(long retryDelayMillis){ 108 queue(this, retryDelayMillis); 109 } 110 /** 111 * Takes care of non transient failures. The job is suspended, and the state of the action is changed to *MANUAL and 112 * set pending flag of action to false 113 * 114 * @param context the execution context. 115 * @param executor the executor instance being used. 116 * @param status the status to be set for the action. 117 * @throws CommandException thrown if unable to suspend job 118 */ 119 protected void handleNonTransient(ActionExecutor.Context context, ActionExecutor executor, 120 WorkflowAction.Status status) throws CommandException { 121 ActionExecutorContext aContext = (ActionExecutorContext) context; 122 WorkflowActionBean action = (WorkflowActionBean) aContext.getAction(); 123 incrActionErrorCounter(action.getType(), "nontransient", 1); 124 WorkflowJobBean workflow = (WorkflowJobBean) context.getWorkflow(); 125 String id = workflow.getId(); 126 action.setStatus(status); 127 action.resetPendingOnly(); 128 LOG.warn("Suspending Workflow Job id=" + id); 129 try { 130 SuspendXCommand.suspendJob(Services.get().get(JPAService.class), workflow, id, action.getId(), null); 131 } 132 catch (Exception e) { 133 throw new CommandException(ErrorCode.E0727, id, e.getMessage()); 134 } 135 finally { 136 updateParentIfNecessary(workflow, 3); 137 } 138 } 139 140 /** 141 * Takes care of errors. </p> For errors while attempting to start the action, the job state is updated and an 142 * {@link ActionEndCommand} is queued. </p> For errors while attempting to end the action, the job state is updated. 143 * </p> 144 * 145 * @param context the execution context. 146 * @param executor the executor instance being used. 147 * @param message 148 * @param isStart whether the error was generated while starting or ending an action. 149 * @param status the status to be set for the action. 150 * @throws CommandException thrown if unable to handle action error 151 */ 152 protected void handleError(ActionExecutor.Context context, ActionExecutor executor, String message, 153 boolean isStart, WorkflowAction.Status status) throws CommandException { 154 LOG.warn("Setting Action Status to [{0}]", status); 155 ActionExecutorContext aContext = (ActionExecutorContext) context; 156 WorkflowActionBean action = (WorkflowActionBean) aContext.getAction(); 157 158 if (!handleUserRetry(action)) { 159 incrActionErrorCounter(action.getType(), "error", 1); 160 action.setPending(); 161 if (isStart) { 162 action.setExecutionData(message, null); 163 queue(new ActionEndXCommand(action.getId(), action.getType())); 164 } 165 else { 166 action.setEndData(status, WorkflowAction.Status.ERROR.toString()); 167 } 168 } 169 } 170 171 /** 172 * Fail the job due to failed action 173 * 174 * @param context the execution context. 175 * @throws CommandException thrown if unable to fail job 176 */ 177 public void failJob(ActionExecutor.Context context) throws CommandException { 178 ActionExecutorContext aContext = (ActionExecutorContext) context; 179 WorkflowActionBean action = (WorkflowActionBean) aContext.getAction(); 180 failJob(context, action); 181 } 182 183 /** 184 * Fail the job due to failed action 185 * 186 * @param context the execution context. 187 * @param action the action that caused the workflow to fail 188 * @throws CommandException thrown if unable to fail job 189 */ 190 public void failJob(ActionExecutor.Context context, WorkflowActionBean action) throws CommandException { 191 WorkflowJobBean workflow = (WorkflowJobBean) context.getWorkflow(); 192 if (!handleUserRetry(action)) { 193 incrActionErrorCounter(action.getType(), "failed", 1); 194 LOG.warn("Failing Job due to failed action [{0}]", action.getName()); 195 try { 196 workflow.getWorkflowInstance().fail(action.getName()); 197 WorkflowInstance wfInstance = workflow.getWorkflowInstance(); 198 ((LiteWorkflowInstance) wfInstance).setStatus(WorkflowInstance.Status.FAILED); 199 workflow.setWorkflowInstance(wfInstance); 200 workflow.setStatus(WorkflowJob.Status.FAILED); 201 action.setStatus(WorkflowAction.Status.FAILED); 202 action.resetPending(); 203 queue(new WorkflowNotificationXCommand(workflow, action)); 204 queue(new KillXCommand(workflow.getId())); 205 InstrumentUtils.incrJobCounter(INSTR_FAILED_JOBS_COUNTER_NAME, 1, getInstrumentation()); 206 } 207 catch (WorkflowException ex) { 208 throw new CommandException(ex); 209 } 210 } 211 } 212 213 /** 214 * Execute retry for action if this action is eligible for user-retry 215 * 216 * @param context the execution context. 217 * @return true if user-retry has to be handled for this action 218 * @throws CommandException thrown if unable to fail job 219 */ 220 public boolean handleUserRetry(WorkflowActionBean action) throws CommandException { 221 String errorCode = action.getErrorCode(); 222 Set<String> allowedRetryCode = LiteWorkflowStoreService.getUserRetryErrorCode(); 223 224 if ((allowedRetryCode.contains(LiteWorkflowStoreService.USER_ERROR_CODE_ALL) || allowedRetryCode.contains(errorCode)) 225 && action.getUserRetryCount() < action.getUserRetryMax()) { 226 LOG.info("Preparing retry this action [{0}], errorCode [{1}], userRetryCount [{2}], " 227 + "userRetryMax [{3}], userRetryInterval [{4}]", action.getId(), errorCode, action 228 .getUserRetryCount(), action.getUserRetryMax(), action.getUserRetryInterval()); 229 int interval = action.getUserRetryInterval() * 60 * 1000; 230 action.setStatus(WorkflowAction.Status.USER_RETRY); 231 action.incrmentUserRetryCount(); 232 action.setPending(); 233 queue(new ActionStartXCommand(action.getId(), action.getType()), interval); 234 return true; 235 } 236 return false; 237 } 238 239 /* 240 * In case of action error increment the error count for instrumentation 241 */ 242 private void incrActionErrorCounter(String type, String error, int count) { 243 getInstrumentation().incr(INSTRUMENTATION_GROUP, type + "#ex." + error, count); 244 } 245 246 /** 247 * Increment the action counter in the instrumentation log. indicating how 248 * many times the action was executed since the start Oozie server 249 */ 250 protected void incrActionCounter(String type, int count) { 251 getInstrumentation().incr(INSTRUMENTATION_GROUP, type + "#" + getName(), count); 252 } 253 254 /** 255 * Adding a cron for the instrumentation time for the given Instrumentation 256 * group 257 */ 258 protected void addActionCron(String type, Instrumentation.Cron cron) { 259 getInstrumentation().addCron(INSTRUMENTATION_GROUP, type + "#" + getName(), cron); 260 } 261 262 /* 263 * Returns the next retry time in milliseconds, based on retry policy algorithm. 264 */ 265 private long getRetryDelay(int retryCount, long retryInterval, ActionExecutor.RETRYPOLICY retryPolicy) { 266 switch (retryPolicy) { 267 case EXPONENTIAL: 268 long retryTime = ((long) Math.pow(2, retryCount) * retryInterval * 1000L); 269 return retryTime; 270 case PERIODIC: 271 return retryInterval * 1000L; 272 default: 273 throw new UnsupportedOperationException("Retry policy not supported"); 274 } 275 } 276 277 /** 278 * Workflow action executor context 279 * 280 */ 281 public static class ActionExecutorContext implements ActionExecutor.Context { 282 protected final WorkflowJobBean workflow; 283 private Configuration protoConf; 284 protected final WorkflowActionBean action; 285 private final boolean isRetry; 286 private final boolean isUserRetry; 287 private boolean started; 288 private boolean ended; 289 private boolean executed; 290 private boolean shouldEndWF; 291 private Job.Status jobStatus; 292 293 /** 294 * Constructing the ActionExecutorContext, setting the private members 295 * and constructing the proto configuration 296 */ 297 public ActionExecutorContext(WorkflowJobBean workflow, WorkflowActionBean action, boolean isRetry, boolean isUserRetry) { 298 this.workflow = workflow; 299 this.action = action; 300 this.isRetry = isRetry; 301 this.isUserRetry = isUserRetry; 302 try { 303 protoConf = new XConfiguration(new StringReader(workflow.getProtoActionConf())); 304 } 305 catch (IOException ex) { 306 throw new RuntimeException("It should not happen", ex); 307 } 308 } 309 310 /* 311 * (non-Javadoc) 312 * @see org.apache.oozie.action.ActionExecutor.Context#getCallbackUrl(java.lang.String) 313 */ 314 public String getCallbackUrl(String externalStatusVar) { 315 return Services.get().get(CallbackService.class).createCallBackUrl(action.getId(), externalStatusVar); 316 } 317 318 /* 319 * (non-Javadoc) 320 * @see org.apache.oozie.action.ActionExecutor.Context#getProtoActionConf() 321 */ 322 public Configuration getProtoActionConf() { 323 return protoConf; 324 } 325 326 /* 327 * (non-Javadoc) 328 * @see org.apache.oozie.action.ActionExecutor.Context#getWorkflow() 329 */ 330 public WorkflowJob getWorkflow() { 331 return workflow; 332 } 333 334 /** 335 * Returns the workflow action of the given action context 336 * 337 * @return the workflow action of the given action context 338 */ 339 public WorkflowAction getAction() { 340 return action; 341 } 342 343 /* 344 * (non-Javadoc) 345 * @see org.apache.oozie.action.ActionExecutor.Context#getELEvaluator() 346 */ 347 public ELEvaluator getELEvaluator() { 348 ELEvaluator evaluator = Services.get().get(ELService.class).createEvaluator("workflow"); 349 DagELFunctions.configureEvaluator(evaluator, workflow, action); 350 return evaluator; 351 } 352 353 /* 354 * (non-Javadoc) 355 * @see org.apache.oozie.action.ActionExecutor.Context#setVar(java.lang.String, java.lang.String) 356 */ 357 public void setVar(String name, String value) { 358 setVarToWorkflow(name, value); 359 } 360 361 /** 362 * This is not thread safe, don't use if workflowjob is shared among multiple actions command 363 * @param name 364 * @param value 365 */ 366 public void setVarToWorkflow(String name, String value) { 367 name = action.getName() + WorkflowInstance.NODE_VAR_SEPARATOR + name; 368 WorkflowInstance wfInstance = workflow.getWorkflowInstance(); 369 wfInstance.setVar(name, value); 370 workflow.setWorkflowInstance(wfInstance); 371 } 372 373 /* 374 * (non-Javadoc) 375 * @see org.apache.oozie.action.ActionExecutor.Context#getVar(java.lang.String) 376 */ 377 public String getVar(String name) { 378 name = action.getName() + WorkflowInstance.NODE_VAR_SEPARATOR + name; 379 return workflow.getWorkflowInstance().getVar(name); 380 } 381 382 /* 383 * (non-Javadoc) 384 * @see org.apache.oozie.action.ActionExecutor.Context#setStartData(java.lang.String, java.lang.String, java.lang.String) 385 */ 386 public void setStartData(String externalId, String trackerUri, String consoleUrl) { 387 action.setStartData(externalId, trackerUri, consoleUrl); 388 started = true; 389 } 390 391 /** 392 * Setting the start time of the action 393 */ 394 public void setStartTime() { 395 Date now = new Date(); 396 action.setStartTime(now); 397 } 398 399 /* 400 * (non-Javadoc) 401 * @see org.apache.oozie.action.ActionExecutor.Context#setExecutionData(java.lang.String, java.util.Properties) 402 */ 403 public void setExecutionData(String externalStatus, Properties actionData) { 404 action.setExecutionData(externalStatus, actionData); 405 executed = true; 406 } 407 408 /* 409 * (non-Javadoc) 410 * @see org.apache.oozie.action.ActionExecutor.Context#setExecutionStats(java.lang.String) 411 */ 412 public void setExecutionStats(String jsonStats) { 413 action.setExecutionStats(jsonStats); 414 executed = true; 415 } 416 417 /* 418 * (non-Javadoc) 419 * @see org.apache.oozie.action.ActionExecutor.Context#setExternalChildIDs(java.lang.String) 420 */ 421 public void setExternalChildIDs(String externalChildIDs) { 422 action.setExternalChildIDs(externalChildIDs); 423 executed = true; 424 } 425 426 /* 427 * (non-Javadoc) 428 * @see org.apache.oozie.action.ActionExecutor.Context#setEndData(org.apache.oozie.client.WorkflowAction.Status, java.lang.String) 429 */ 430 public void setEndData(WorkflowAction.Status status, String signalValue) { 431 action.setEndData(status, signalValue); 432 ended = true; 433 } 434 435 /* 436 * (non-Javadoc) 437 * @see org.apache.oozie.action.ActionExecutor.Context#isRetry() 438 */ 439 public boolean isRetry() { 440 return isRetry; 441 } 442 443 /** 444 * Return if the executor invocation is a user retry or not. 445 * 446 * @return if the executor invocation is a user retry or not. 447 */ 448 public boolean isUserRetry() { 449 return isUserRetry; 450 } 451 452 /** 453 * Returns whether setStartData has been called or not. 454 * 455 * @return true if start completion info has been set. 456 */ 457 public boolean isStarted() { 458 return started; 459 } 460 461 /** 462 * Returns whether setExecutionData has been called or not. 463 * 464 * @return true if execution completion info has been set, otherwise false. 465 */ 466 public boolean isExecuted() { 467 return executed; 468 } 469 470 /** 471 * Returns whether setEndData has been called or not. 472 * 473 * @return true if end completion info has been set. 474 */ 475 public boolean isEnded() { 476 return ended; 477 } 478 479 public void setExternalStatus(String externalStatus) { 480 action.setExternalStatus(externalStatus); 481 } 482 483 @Override 484 public String getRecoveryId() { 485 return action.getId() + RECOVERY_ID_SEPARATOR + workflow.getRun(); 486 } 487 488 /* (non-Javadoc) 489 * @see org.apache.oozie.action.ActionExecutor.Context#getActionDir() 490 */ 491 public Path getActionDir() throws HadoopAccessorException, IOException, URISyntaxException { 492 String name = getWorkflow().getId() + "/" + action.getName() + "--" + action.getType(); 493 FileSystem fs = getAppFileSystem(); 494 String actionDirPath = Services.get().getSystemId() + "/" + name; 495 Path fqActionDir = new Path(fs.getHomeDirectory(), actionDirPath); 496 return fqActionDir; 497 } 498 499 /* (non-Javadoc) 500 * @see org.apache.oozie.action.ActionExecutor.Context#getAppFileSystem() 501 */ 502 public FileSystem getAppFileSystem() throws HadoopAccessorException, IOException, URISyntaxException { 503 WorkflowJob workflow = getWorkflow(); 504 URI uri = new URI(getWorkflow().getAppPath()); 505 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 506 Configuration fsConf = has.createJobConf(uri.getAuthority()); 507 return has.createFileSystem(workflow.getUser(), uri, fsConf); 508 509 } 510 511 /* (non-Javadoc) 512 * @see org.apache.oozie.action.ActionExecutor.Context#setErrorInfo(java.lang.String, java.lang.String) 513 */ 514 @Override 515 public void setErrorInfo(String str, String exMsg) { 516 action.setErrorInfo(str, exMsg); 517 } 518 519 public boolean isShouldEndWF() { 520 return shouldEndWF; 521 } 522 523 public void setShouldEndWF(boolean shouldEndWF) { 524 this.shouldEndWF = shouldEndWF; 525 } 526 527 public Job.Status getJobStatus() { 528 return jobStatus; 529 } 530 531 public void setJobStatus(Job.Status jobStatus) { 532 this.jobStatus = jobStatus; 533 } 534 } 535 536 public static class ForkedActionExecutorContext extends ActionExecutorContext { 537 private Map<String, String> contextVariableMap = new HashMap<String, String>(); 538 539 public ForkedActionExecutorContext(WorkflowJobBean workflow, WorkflowActionBean action, boolean isRetry, 540 boolean isUserRetry) { 541 super(workflow, action, isRetry, isUserRetry); 542 } 543 544 public void setVar(String name, String value) { 545 if (value != null) { 546 contextVariableMap.remove(name); 547 } 548 else { 549 contextVariableMap.put(name, value); 550 } 551 } 552 553 public String getVar(String name) { 554 return contextVariableMap.get(name); 555 } 556 557 public Map<String, String> getContextMap() { 558 return contextVariableMap; 559 } 560 } 561 562}