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.util.ArrayList; 022import java.util.Date; 023import java.util.List; 024import javax.servlet.jsp.el.ELException; 025 026import org.apache.hadoop.conf.Configuration; 027import org.apache.oozie.ErrorCode; 028import org.apache.oozie.FaultInjection; 029import org.apache.oozie.SLAEventBean; 030import org.apache.oozie.WorkflowActionBean; 031import org.apache.oozie.WorkflowJobBean; 032import org.apache.oozie.XException; 033import org.apache.oozie.action.ActionExecutor; 034import org.apache.oozie.action.ActionExecutorException; 035import org.apache.oozie.action.control.ControlNodeActionExecutor; 036import org.apache.oozie.client.OozieClient; 037import org.apache.oozie.client.WorkflowAction; 038import org.apache.oozie.client.WorkflowJob; 039import org.apache.oozie.client.SLAEvent.SlaAppType; 040import org.apache.oozie.client.SLAEvent.Status; 041import org.apache.oozie.client.rest.JsonBean; 042import org.apache.oozie.command.CommandException; 043import org.apache.oozie.command.PreconditionException; 044import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry; 045import org.apache.oozie.executor.jpa.BatchQueryExecutor; 046import org.apache.oozie.executor.jpa.JPAExecutorException; 047import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor; 048import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor.WorkflowActionQuery; 049import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor; 050import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor.WorkflowJobQuery; 051import org.apache.oozie.service.ActionService; 052import org.apache.oozie.service.EventHandlerService; 053import org.apache.oozie.service.JPAService; 054import org.apache.oozie.service.Services; 055import org.apache.oozie.service.UUIDService; 056import org.apache.oozie.util.ELEvaluationException; 057import org.apache.oozie.util.Instrumentation; 058import org.apache.oozie.util.LogUtils; 059import org.apache.oozie.util.XLog; 060import org.apache.oozie.util.XmlUtils; 061import org.apache.oozie.util.db.SLADbXOperations; 062 063@SuppressWarnings("deprecation") 064public class ActionStartXCommand extends ActionXCommand<org.apache.oozie.command.wf.ActionXCommand.ActionExecutorContext> { 065 public static final String EL_ERROR = "EL_ERROR"; 066 public static final String EL_EVAL_ERROR = "EL_EVAL_ERROR"; 067 public static final String COULD_NOT_START = "COULD_NOT_START"; 068 public static final String START_DATA_MISSING = "START_DATA_MISSING"; 069 public static final String EXEC_DATA_MISSING = "EXEC_DATA_MISSING"; 070 071 private String jobId = null; 072 protected String actionId = null; 073 protected WorkflowJobBean wfJob = null; 074 protected WorkflowActionBean wfAction = null; 075 private JPAService jpaService = null; 076 private ActionExecutor executor = null; 077 private List<UpdateEntry> updateList = new ArrayList<UpdateEntry>(); 078 private List<JsonBean> insertList = new ArrayList<JsonBean>(); 079 protected ActionExecutorContext context = null; 080 081 public ActionStartXCommand(String actionId, String type) { 082 super("action.start", type, 0); 083 this.actionId = actionId; 084 this.jobId = Services.get().get(UUIDService.class).getId(actionId); 085 } 086 087 public ActionStartXCommand(WorkflowJobBean job, String actionId, String type) { 088 super("action.start", type, 0); 089 this.actionId = actionId; 090 this.wfJob = job; 091 this.jobId = wfJob.getId(); 092 } 093 094 @Override 095 protected void setLogInfo() { 096 LogUtils.setLogInfo(actionId); 097 } 098 099 @Override 100 protected boolean isLockRequired() { 101 return true; 102 } 103 104 @Override 105 public String getEntityKey() { 106 return this.jobId; 107 } 108 109 @Override 110 protected void loadState() throws CommandException { 111 try { 112 jpaService = Services.get().get(JPAService.class); 113 if (jpaService != null) { 114 if (wfJob == null) { 115 this.wfJob = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW, jobId); 116 } 117 this.wfAction = WorkflowActionQueryExecutor.getInstance().get(WorkflowActionQuery.GET_ACTION, actionId); 118 LogUtils.setLogInfo( wfJob); 119 LogUtils.setLogInfo(wfAction); 120 } 121 else { 122 throw new CommandException(ErrorCode.E0610); 123 } 124 } 125 catch (XException ex) { 126 throw new CommandException(ex); 127 } 128 } 129 130 @Override 131 protected void verifyPrecondition() throws CommandException, PreconditionException { 132 if (wfJob == null) { 133 throw new PreconditionException(ErrorCode.E0604, jobId); 134 } 135 if (wfAction == null) { 136 throw new PreconditionException(ErrorCode.E0605, actionId); 137 } 138 if (wfAction.isPending() 139 && (wfAction.getStatus() == WorkflowActionBean.Status.PREP 140 || wfAction.getStatus() == WorkflowActionBean.Status.START_RETRY 141 || wfAction.getStatus() == WorkflowActionBean.Status.START_MANUAL 142 || wfAction.getStatus() == WorkflowActionBean.Status.USER_RETRY 143 )) { 144 if (wfJob.getStatus() != WorkflowJob.Status.RUNNING) { 145 throw new PreconditionException(ErrorCode.E0810, WorkflowJob.Status.RUNNING.toString()); 146 } 147 } 148 else { 149 throw new PreconditionException(ErrorCode.E0816, wfAction.isPending(), wfAction.getStatusStr()); 150 } 151 152 executor = Services.get().get(ActionService.class).getExecutor(wfAction.getType()); 153 if (executor == null) { 154 throw new CommandException(ErrorCode.E0802, wfAction.getType()); 155 } 156 } 157 158 @Override 159 protected ActionExecutorContext execute() throws CommandException { 160 LOG.debug("STARTED ActionStartXCommand for wf actionId=" + actionId); 161 Configuration conf = wfJob.getWorkflowInstance().getConf(); 162 163 int maxRetries = 0; 164 long retryInterval = 0; 165 boolean execSynchronous = false; 166 167 if (!(executor instanceof ControlNodeActionExecutor)) { 168 maxRetries = conf.getInt(OozieClient.ACTION_MAX_RETRIES, executor.getMaxRetries()); 169 retryInterval = conf.getLong(OozieClient.ACTION_RETRY_INTERVAL, executor.getRetryInterval()); 170 } 171 172 executor.setMaxRetries(maxRetries); 173 executor.setRetryInterval(retryInterval); 174 175 try { 176 boolean isRetry = false; 177 if (wfAction.getStatus() == WorkflowActionBean.Status.START_RETRY 178 || wfAction.getStatus() == WorkflowActionBean.Status.START_MANUAL) { 179 isRetry = true; 180 prepareForRetry(wfAction); 181 } 182 boolean isUserRetry = false; 183 if (wfAction.getStatus() == WorkflowActionBean.Status.USER_RETRY) { 184 isUserRetry = true; 185 prepareForRetry(wfAction); 186 } 187 context = getContext(isRetry, isUserRetry); 188 boolean caught = false; 189 try { 190 if (!(executor instanceof ControlNodeActionExecutor)) { 191 String tmpActionConf = XmlUtils.removeComments(wfAction.getConf()); 192 String actionConf = context.getELEvaluator().evaluate(tmpActionConf, String.class); 193 wfAction.setConf(actionConf); 194 LOG.debug("Start, name [{0}] type [{1}] configuration{E}{E}{2}{E}", wfAction.getName(), wfAction 195 .getType(), actionConf); 196 } 197 } 198 catch (ELEvaluationException ex) { 199 caught = true; 200 throw new ActionExecutorException(ActionExecutorException.ErrorType.TRANSIENT, EL_EVAL_ERROR, ex 201 .getMessage(), ex); 202 } 203 catch (ELException ex) { 204 caught = true; 205 context.setErrorInfo(EL_ERROR, ex.getMessage()); 206 LOG.warn("ELException in ActionStartXCommand ", ex.getMessage(), ex); 207 handleError(context, wfJob, wfAction); 208 } 209 catch (org.jdom.JDOMException je) { 210 caught = true; 211 context.setErrorInfo("ParsingError", je.getMessage()); 212 LOG.warn("JDOMException in ActionStartXCommand ", je.getMessage(), je); 213 handleError(context, wfJob, wfAction); 214 } 215 catch (Exception ex) { 216 caught = true; 217 context.setErrorInfo(EL_ERROR, ex.getMessage()); 218 LOG.warn("Exception in ActionStartXCommand ", ex.getMessage(), ex); 219 handleError(context, wfJob, wfAction); 220 } 221 if(!caught) { 222 wfAction.setErrorInfo(null, null); 223 incrActionCounter(wfAction.getType(), 1); 224 225 LOG.info("Start action [{0}] with user-retry state : userRetryCount [{1}], userRetryMax [{2}], userRetryInterval [{3}]", 226 wfAction.getId(), wfAction.getUserRetryCount(), wfAction.getUserRetryMax(), wfAction 227 .getUserRetryInterval()); 228 229 Instrumentation.Cron cron = new Instrumentation.Cron(); 230 cron.start(); 231 context.setStartTime(); 232 executor.start(context, wfAction); 233 cron.stop(); 234 FaultInjection.activate("org.apache.oozie.command.SkipCommitFaultInjection"); 235 addActionCron(wfAction.getType(), cron); 236 237 wfAction.setRetries(0); 238 if (wfAction.isExecutionComplete()) { 239 if (!context.isExecuted()) { 240 LOG.warn(XLog.OPS, "Action Completed, ActionExecutor [{0}] must call setExecutionData()", executor 241 .getType()); 242 wfAction.setErrorInfo(EXEC_DATA_MISSING, 243 "Execution Complete, but Execution Data Missing from Action"); 244 failJob(context); 245 } else { 246 wfAction.setPending(); 247 if (!(executor instanceof ControlNodeActionExecutor)) { 248 queue(new ActionEndXCommand(wfAction.getId(), wfAction.getType())); 249 } 250 else { 251 execSynchronous = true; 252 } 253 } 254 } 255 else { 256 if (!context.isStarted()) { 257 LOG.warn(XLog.OPS, "Action Started, ActionExecutor [{0}] must call setStartData()", executor 258 .getType()); 259 wfAction.setErrorInfo(START_DATA_MISSING, "Execution Started, but Start Data Missing from Action"); 260 failJob(context); 261 } else { 262 queue(new WorkflowNotificationXCommand(wfJob, wfAction)); 263 } 264 } 265 266 LOG.info(XLog.STD, "[***" + wfAction.getId() + "***]" + "Action status=" + wfAction.getStatusStr()); 267 268 updateList.add(new UpdateEntry<WorkflowActionQuery>(WorkflowActionQuery.UPDATE_ACTION_START, wfAction)); 269 updateJobLastModified(); 270 // Add SLA status event (STARTED) for WF_ACTION 271 SLAEventBean slaEvent = SLADbXOperations.createStatusEvent(wfAction.getSlaXml(), wfAction.getId(), Status.STARTED, 272 SlaAppType.WORKFLOW_ACTION); 273 if(slaEvent != null) { 274 insertList.add(slaEvent); 275 } 276 LOG.info(XLog.STD, "[***" + wfAction.getId() + "***]" + "Action updated in DB!"); 277 } 278 } 279 catch (ActionExecutorException ex) { 280 LOG.warn("Error starting action [{0}]. ErrorType [{1}], ErrorCode [{2}], Message [{3}]", 281 wfAction.getName(), ex.getErrorType(), ex.getErrorCode(), ex.getMessage(), ex); 282 wfAction.setErrorInfo(ex.getErrorCode(), ex.getMessage()); 283 switch (ex.getErrorType()) { 284 case TRANSIENT: 285 if (!handleTransient(context, executor, WorkflowAction.Status.START_RETRY)) { 286 handleNonTransient(context, executor, WorkflowAction.Status.START_MANUAL); 287 wfAction.setPendingAge(new Date()); 288 wfAction.setRetries(0); 289 wfAction.setStartTime(null); 290 } 291 break; 292 case NON_TRANSIENT: 293 handleNonTransient(context, executor, WorkflowAction.Status.START_MANUAL); 294 break; 295 case ERROR: 296 handleError(context, executor, WorkflowAction.Status.ERROR.toString(), true, 297 WorkflowAction.Status.DONE); 298 break; 299 case FAILED: 300 try { 301 failJob(context); 302 endWF(); 303 SLAEventBean slaEvent1 = SLADbXOperations.createStatusEvent(wfAction.getSlaXml(), wfAction.getId(), Status.FAILED, 304 SlaAppType.WORKFLOW_ACTION); 305 if(slaEvent1 != null) { 306 insertList.add(slaEvent1); 307 } 308 } 309 catch (XException x) { 310 LOG.warn("ActionStartXCommand - case:FAILED ", x.getMessage()); 311 } 312 break; 313 } 314 updateList.add(new UpdateEntry<WorkflowActionQuery>(WorkflowActionQuery.UPDATE_ACTION_START, wfAction)); 315 updateJobLastModified(); 316 } 317 finally { 318 try { 319 BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, updateList, null); 320 if (!(executor instanceof ControlNodeActionExecutor) && EventHandlerService.isEnabled()) { 321 generateEvent(wfAction, wfJob.getUser()); 322 } 323 if (execSynchronous) { 324 // Changing to synchronous call from asynchronous queuing to prevent 325 // undue delay from ::start:: to action due to queuing 326 callActionEnd(); 327 } 328 } 329 catch (JPAExecutorException e) { 330 throw new CommandException(e); 331 } 332 } 333 334 LOG.debug("ENDED ActionStartXCommand for wf actionId=" + actionId + ", jobId=" + jobId); 335 336 return null; 337 } 338 339 protected void callActionEnd() throws CommandException { 340 new ActionEndXCommand(wfAction.getId(), wfAction.getType()).call(getEntityKey()); 341 } 342 343 /** 344 * Get action executor context 345 * @param isRetry 346 * @param isUserRetry 347 * @return 348 */ 349 protected ActionExecutorContext getContext(boolean isRetry, boolean isUserRetry) { 350 return new ActionXCommand.ActionExecutorContext(wfJob, wfAction, isRetry, isUserRetry); 351 } 352 353 protected void updateJobLastModified(){ 354 wfJob.setLastModifiedTime(new Date()); 355 updateList.add(new UpdateEntry<WorkflowJobQuery>(WorkflowJobQuery.UPDATE_WORKFLOW_STATUS_INSTANCE_MODIFIED, wfJob)); 356 } 357 358 protected void endWF() throws CommandException{ 359 updateParentIfNecessary(wfJob, 3); 360 new WfEndXCommand(wfJob).call(); // To delete the WF temp dir 361 SLAEventBean slaEvent2 = SLADbXOperations.createStatusEvent(wfJob.getSlaXml(), wfJob.getId(), Status.FAILED, 362 SlaAppType.WORKFLOW_JOB); 363 if(slaEvent2 != null) { 364 insertList.add(slaEvent2); 365 } 366 } 367 368 protected void handleError(ActionExecutorContext context, WorkflowJobBean workflow, WorkflowActionBean action) 369 throws CommandException { 370 failJob(context); 371 updateList.add(new UpdateEntry<WorkflowActionQuery>(WorkflowActionQuery.UPDATE_ACTION_START, wfAction)); 372 updateJobLastModified(); 373 SLAEventBean slaEvent1 = SLADbXOperations.createStatusEvent(action.getSlaXml(), action.getId(), 374 Status.FAILED, SlaAppType.WORKFLOW_ACTION); 375 if(slaEvent1 != null) { 376 insertList.add(slaEvent1); 377 } 378 endWF(); 379 return; 380 } 381 382 /* (non-Javadoc) 383 * @see org.apache.oozie.command.XCommand#getKey() 384 */ 385 @Override 386 public String getKey(){ 387 return getName() + "_" + actionId; 388 } 389 390 private void prepareForRetry(WorkflowActionBean wfAction) { 391 if (wfAction.getType().equals("map-reduce")) { 392 // need to delete child job id of original run 393 wfAction.setExternalChildIDs(""); 394 } 395 } 396 397 @Override 398 protected void queueCommandForTransientFailure(long retryDelayMillis){ 399 queue(new ActionStartXCommand(wfAction.getId(), wfAction.getType()), retryDelayMillis); 400 } 401 402}