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.coord; 020 021import org.apache.hadoop.conf.Configuration; 022import org.apache.oozie.AppType; 023import org.apache.oozie.CoordinatorActionBean; 024import org.apache.oozie.CoordinatorJobBean; 025import org.apache.oozie.ErrorCode; 026import org.apache.oozie.SLAEventBean; 027import org.apache.oozie.client.CoordinatorJob; 028import org.apache.oozie.client.Job; 029import org.apache.oozie.client.SLAEvent.SlaAppType; 030import org.apache.oozie.client.rest.JsonBean; 031import org.apache.oozie.command.CommandException; 032import org.apache.oozie.command.MaterializeTransitionXCommand; 033import org.apache.oozie.command.PreconditionException; 034import org.apache.oozie.command.bundle.BundleStatusUpdateXCommand; 035import org.apache.oozie.coord.TimeUnit; 036import org.apache.oozie.executor.jpa.BatchQueryExecutor; 037import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry; 038import org.apache.oozie.executor.jpa.CoordActionsActiveCountJPAExecutor; 039import org.apache.oozie.executor.jpa.CoordJobQueryExecutor; 040import org.apache.oozie.executor.jpa.CoordJobQueryExecutor.CoordJobQuery; 041import org.apache.oozie.executor.jpa.JPAExecutorException; 042import org.apache.oozie.service.ConfigurationService; 043import org.apache.oozie.service.CoordMaterializeTriggerService; 044import org.apache.oozie.service.EventHandlerService; 045import org.apache.oozie.service.JPAService; 046import org.apache.oozie.service.Service; 047import org.apache.oozie.service.Services; 048import org.apache.oozie.sla.SLAOperations; 049import org.apache.oozie.util.DateUtils; 050import org.apache.oozie.util.Instrumentation; 051import org.apache.oozie.util.LogUtils; 052import org.apache.oozie.util.ParamChecker; 053import org.apache.oozie.util.StatusUtils; 054import org.apache.oozie.util.XConfiguration; 055import org.apache.oozie.util.XmlUtils; 056import org.apache.oozie.util.db.SLADbOperations; 057import org.jdom.Element; 058import org.jdom.JDOMException; 059 060import java.io.IOException; 061import java.io.StringReader; 062import java.sql.Timestamp; 063import java.util.Calendar; 064import java.util.Date; 065import java.util.TimeZone; 066 067/** 068 * Materialize actions for specified start and end time for coordinator job. 069 */ 070@SuppressWarnings("deprecation") 071public class CoordMaterializeTransitionXCommand extends MaterializeTransitionXCommand { 072 073 private JPAService jpaService = null; 074 private CoordinatorJobBean coordJob = null; 075 private String jobId = null; 076 private Date startMatdTime = null; 077 private Date endMatdTime = null; 078 private final int materializationWindow; 079 private int lastActionNumber = 1; // over-ride by DB value 080 private CoordinatorJob.Status prevStatus = null; 081 082 static final private int lookAheadWindow = ConfigurationService.getInt(CoordMaterializeTriggerService 083 .CONF_LOOKUP_INTERVAL); 084 085 /** 086 * Default MAX timeout in minutes, after which coordinator input check will timeout 087 */ 088 public static final String CONF_DEFAULT_MAX_TIMEOUT = Service.CONF_PREFIX + "coord.default.max.timeout"; 089 090 /** 091 * The constructor for class {@link CoordMaterializeTransitionXCommand} 092 * 093 * @param jobId coordinator job id 094 * @param materializationWindow materialization window to calculate end time 095 */ 096 public CoordMaterializeTransitionXCommand(String jobId, int materializationWindow) { 097 super("coord_mater", "coord_mater", 1); 098 this.jobId = ParamChecker.notEmpty(jobId, "jobId"); 099 this.materializationWindow = materializationWindow; 100 } 101 102 public CoordMaterializeTransitionXCommand(CoordinatorJobBean coordJob, int materializationWindow, Date startTime, 103 Date endTime) { 104 super("coord_mater", "coord_mater", 1); 105 this.jobId = ParamChecker.notEmpty(coordJob.getId(), "jobId"); 106 this.materializationWindow = materializationWindow; 107 this.coordJob = coordJob; 108 this.startMatdTime = startTime; 109 this.endMatdTime = endTime; 110 } 111 112 /* (non-Javadoc) 113 * @see org.apache.oozie.command.MaterializeTransitionXCommand#transitToNext() 114 */ 115 @Override 116 public void transitToNext() throws CommandException { 117 } 118 119 /* (non-Javadoc) 120 * @see org.apache.oozie.command.TransitionXCommand#updateJob() 121 */ 122 @Override 123 public void updateJob() throws CommandException { 124 updateList.add(new UpdateEntry(CoordJobQuery.UPDATE_COORD_JOB_MATERIALIZE,coordJob)); 125 } 126 127 /* (non-Javadoc) 128 * @see org.apache.oozie.command.MaterializeTransitionXCommand#performWrites() 129 */ 130 @Override 131 public void performWrites() throws CommandException { 132 try { 133 BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, updateList, null); 134 // register the partition related dependencies of actions 135 for (JsonBean actionBean : insertList) { 136 if (actionBean instanceof CoordinatorActionBean) { 137 CoordinatorActionBean coordAction = (CoordinatorActionBean) actionBean; 138 if (EventHandlerService.isEnabled()) { 139 CoordinatorXCommand.generateEvent(coordAction, coordJob.getUser(), coordJob.getAppName(), null); 140 } 141 142 // TODO: time 100s should be configurable 143 queue(new CoordActionNotificationXCommand(coordAction), 100); 144 145 //Delay for input check = (nominal time - now) 146 long checkDelay = coordAction.getNominalTime().getTime() - new Date().getTime(); 147 queue(new CoordActionInputCheckXCommand(coordAction.getId(), coordAction.getJobId()), 148 Math.max(checkDelay, 0)); 149 150 if (coordAction.getPushMissingDependencies() != null) { 151 // TODO: Delay in catchup mode? 152 queue(new CoordPushDependencyCheckXCommand(coordAction.getId(), true), 100); 153 } 154 } 155 } 156 } 157 catch (JPAExecutorException jex) { 158 throw new CommandException(jex); 159 } 160 } 161 162 /* (non-Javadoc) 163 * @see org.apache.oozie.command.XCommand#getEntityKey() 164 */ 165 @Override 166 public String getEntityKey() { 167 return this.jobId; 168 } 169 170 @Override 171 protected boolean isLockRequired() { 172 return true; 173 } 174 175 /* (non-Javadoc) 176 * @see org.apache.oozie.command.XCommand#loadState() 177 */ 178 @Override 179 protected void loadState() throws CommandException { 180 jpaService = Services.get().get(JPAService.class); 181 if (jpaService == null) { 182 LOG.error(ErrorCode.E0610); 183 } 184 185 try { 186 coordJob = CoordJobQueryExecutor.getInstance().get(CoordJobQuery.GET_COORD_JOB_MATERIALIZE, jobId); 187 prevStatus = coordJob.getStatus(); 188 } 189 catch (JPAExecutorException jex) { 190 throw new CommandException(jex); 191 } 192 193 // calculate start materialize and end materialize time 194 calcMatdTime(); 195 196 LogUtils.setLogInfo(coordJob); 197 } 198 199 /** 200 * Calculate startMatdTime and endMatdTime from job's start time if next materialized time is null 201 * 202 * @throws CommandException thrown if failed to calculate startMatdTime and endMatdTime 203 */ 204 protected void calcMatdTime() throws CommandException { 205 Timestamp startTime = coordJob.getNextMaterializedTimestamp(); 206 if (startTime == null) { 207 startTime = coordJob.getStartTimestamp(); 208 } 209 // calculate end time by adding materializationWindow to start time. 210 // need to convert materializationWindow from secs to milliseconds 211 long startTimeMilli = startTime.getTime(); 212 long endTimeMilli = startTimeMilli + (materializationWindow * 1000); 213 214 startMatdTime = DateUtils.toDate(new Timestamp(startTimeMilli)); 215 endMatdTime = DateUtils.toDate(new Timestamp(endTimeMilli)); 216 endMatdTime = getMaterializationTimeForCatchUp(endMatdTime); 217 // if MaterializationWindow end time is greater than endTime 218 // for job, then set it to endTime of job 219 Date jobEndTime = coordJob.getEndTime(); 220 if (endMatdTime.compareTo(jobEndTime) > 0) { 221 endMatdTime = jobEndTime; 222 } 223 224 LOG.debug("Materializing coord job id=" + jobId + ", start=" + DateUtils.formatDateOozieTZ(startMatdTime) + ", end=" + DateUtils.formatDateOozieTZ(endMatdTime) 225 + ", window=" + materializationWindow); 226 } 227 228 /** 229 * Get materialization for window for catch-up jobs. for current jobs,it reruns currentMatdate, For catch-up, end 230 * Mataterilized Time = startMatdTime + MatThrottling * frequency; unless LAST_ONLY execution order is set, in which 231 * case it returns now (to materialize all actions in the past) 232 * 233 * @param currentMatTime 234 * @return 235 * @throws CommandException 236 * @throws JDOMException 237 */ 238 private Date getMaterializationTimeForCatchUp(Date currentMatTime) throws CommandException { 239 if (currentMatTime.after(new Date())) { 240 return currentMatTime; 241 } 242 if (coordJob.getExecutionOrder().equals(CoordinatorJob.Execution.LAST_ONLY) || 243 coordJob.getExecutionOrder().equals(CoordinatorJob.Execution.NONE)) { 244 return new Date(); 245 } 246 final int frequency; 247 try { 248 frequency = Integer.parseInt(coordJob.getFrequency()); 249 } 250 catch (final NumberFormatException e) { 251 // Cron based frequency: catching up at maximum till the coordinator job's end time, 252 // bounded also by the throttle parameter, aka the number of coordinator actions to materialize 253 return coordJob.getEndTime(); 254 } 255 256 TimeZone appTz = DateUtils.getTimeZone(coordJob.getTimeZone()); 257 TimeUnit freqTU = TimeUnit.valueOf(coordJob.getTimeUnitStr()); 258 Calendar startInstance = Calendar.getInstance(appTz); 259 startInstance.setTime(startMatdTime); 260 Calendar endMatInstance = null; 261 Calendar previousInstance = startInstance; 262 for (int i = 1; i <= coordJob.getMatThrottling(); i++) { 263 endMatInstance = (Calendar) startInstance.clone(); 264 endMatInstance.add(freqTU.getCalendarUnit(), i * frequency); 265 if (endMatInstance.getTime().compareTo(new Date()) >= 0) { 266 if (previousInstance.getTime().after(currentMatTime)) { 267 return previousInstance.getTime(); 268 } 269 else { 270 return currentMatTime; 271 } 272 } 273 previousInstance = endMatInstance; 274 } 275 if (endMatInstance == null) { 276 return currentMatTime; 277 } 278 else { 279 return endMatInstance.getTime(); 280 } 281 } 282 283 /* (non-Javadoc) 284 * @see org.apache.oozie.command.XCommand#verifyPrecondition() 285 */ 286 @Override 287 protected void verifyPrecondition() throws CommandException, PreconditionException { 288 if (!(coordJob.getStatus() == CoordinatorJobBean.Status.PREP || coordJob.getStatus() == CoordinatorJobBean.Status.RUNNING 289 || coordJob.getStatus() == CoordinatorJobBean.Status.RUNNINGWITHERROR)) { 290 throw new PreconditionException(ErrorCode.E1100, "CoordMaterializeTransitionXCommand for jobId=" + jobId 291 + " job is not in PREP or RUNNING but in " + coordJob.getStatus()); 292 } 293 294 if (coordJob.isDoneMaterialization()) { 295 throw new PreconditionException(ErrorCode.E1100, "CoordMaterializeTransitionXCommand for jobId =" + jobId 296 + " job is already materialized"); 297 } 298 299 if (coordJob.getNextMaterializedTimestamp() != null 300 && coordJob.getNextMaterializedTimestamp().compareTo(coordJob.getEndTimestamp()) >= 0) { 301 throw new PreconditionException(ErrorCode.E1100, "CoordMaterializeTransitionXCommand for jobId=" + jobId 302 + " job is already materialized"); 303 } 304 305 Timestamp startTime = coordJob.getNextMaterializedTimestamp(); 306 if (startTime == null) { 307 startTime = coordJob.getStartTimestamp(); 308 309 if (startTime.after(new Timestamp(System.currentTimeMillis() + lookAheadWindow * 1000))) { 310 throw new PreconditionException(ErrorCode.E1100, "CoordMaterializeTransitionXCommand for jobId=" 311 + jobId + " job's start time is not reached yet - nothing to materialize"); 312 } 313 } 314 315 if (coordJob.getNextMaterializedTimestamp() != null 316 && coordJob.getNextMaterializedTimestamp().after( 317 new Timestamp(System.currentTimeMillis() + lookAheadWindow * 1000))) { 318 throw new PreconditionException(ErrorCode.E1100, "CoordMaterializeTransitionXCommand for jobId=" + jobId 319 + " Request is for future time. Lookup time is " 320 + new Timestamp(System.currentTimeMillis() + lookAheadWindow * 1000) + " mat time is " 321 + coordJob.getNextMaterializedTimestamp()); 322 } 323 324 if (coordJob.getLastActionTime() != null && coordJob.getLastActionTime().compareTo(coordJob.getEndTime()) >= 0) { 325 throw new PreconditionException(ErrorCode.E1100, "ENDED Coordinator materialization for jobId = " + jobId 326 + ", all actions have been materialized from start time = " + coordJob.getStartTime() 327 + " to end time = " + coordJob.getEndTime() + ", job status = " + coordJob.getStatusStr()); 328 } 329 330 if (coordJob.getLastActionTime() != null && coordJob.getLastActionTime().compareTo(endMatdTime) >= 0) { 331 throw new PreconditionException(ErrorCode.E1100, "ENDED Coordinator materialization for jobId = " + jobId 332 + ", action is *already* materialized for Materialization start time = " + startMatdTime 333 + ", materialization end time = " + endMatdTime + ", job status = " + coordJob.getStatusStr()); 334 } 335 336 if (endMatdTime.after(coordJob.getEndTime())) { 337 throw new PreconditionException(ErrorCode.E1100, "ENDED Coordinator materialization for jobId = " + jobId 338 + " materialization end time = " + endMatdTime + " surpasses coordinator job's end time = " 339 + coordJob.getEndTime() + " job status = " + coordJob.getStatusStr()); 340 } 341 342 if (coordJob.getPauseTime() != null && !startMatdTime.before(coordJob.getPauseTime())) { 343 throw new PreconditionException(ErrorCode.E1100, "ENDED Coordinator materialization for jobId = " + jobId 344 + ", materialization start time = " + startMatdTime 345 + " is after or equal to coordinator job's pause time = " + coordJob.getPauseTime() 346 + ", job status = " + coordJob.getStatusStr()); 347 } 348 349 } 350 351 /* (non-Javadoc) 352 * @see org.apache.oozie.command.MaterializeTransitionXCommand#materialize() 353 */ 354 @Override 355 protected void materialize() throws CommandException { 356 Instrumentation.Cron cron = new Instrumentation.Cron(); 357 cron.start(); 358 try { 359 materializeActions(false); 360 updateJobMaterializeInfo(coordJob); 361 } 362 catch (CommandException ex) { 363 LOG.warn("Exception occurred:" + ex.getMessage() + " Making the job failed ", ex); 364 coordJob.setStatus(Job.Status.FAILED); 365 coordJob.resetPending(); 366 // remove any materialized actions and slaEvents 367 insertList.clear(); 368 } 369 catch (Exception e) { 370 LOG.error("Exception occurred:" + e.getMessage() + " Making the job failed ", e); 371 coordJob.setStatus(Job.Status.FAILED); 372 try { 373 CoordJobQueryExecutor.getInstance().executeUpdate(CoordJobQuery.UPDATE_COORD_JOB_MATERIALIZE, coordJob); 374 } 375 catch (JPAExecutorException jex) { 376 throw new CommandException(ErrorCode.E1011, jex); 377 } 378 throw new CommandException(ErrorCode.E1012, e.getMessage(), e); 379 } finally { 380 cron.stop(); 381 instrumentation.addCron(INSTRUMENTATION_GROUP, getName() + ".materialize", cron); 382 } 383 } 384 385 /** 386 * Create action instances starting from "startMatdTime" to "endMatdTime" and store them into coord action table. 387 * 388 * @param dryrun if this is a dry run 389 * @throws Exception thrown if failed to materialize actions 390 */ 391 protected String materializeActions(boolean dryrun) throws Exception { 392 393 Configuration jobConf = null; 394 try { 395 jobConf = new XConfiguration(new StringReader(coordJob.getConf())); 396 } 397 catch (IOException ioe) { 398 LOG.warn("Configuration parse error. read from DB :" + coordJob.getConf(), ioe); 399 throw new CommandException(ErrorCode.E1005, ioe.getMessage(), ioe); 400 } 401 402 String jobXml = coordJob.getJobXml(); 403 Element eJob = XmlUtils.parseXml(jobXml); 404 TimeZone appTz = DateUtils.getTimeZone(coordJob.getTimeZone()); 405 406 String frequency = coordJob.getFrequency(); 407 TimeUnit freqTU = TimeUnit.valueOf(coordJob.getTimeUnitStr()); 408 TimeUnit endOfFlag = TimeUnit.valueOf(eJob.getAttributeValue("end_of_duration")); 409 Calendar start = Calendar.getInstance(appTz); 410 start.setTime(startMatdTime); 411 DateUtils.moveToEnd(start, endOfFlag); 412 Calendar end = Calendar.getInstance(appTz); 413 end.setTime(endMatdTime); 414 lastActionNumber = coordJob.getLastActionNumber(); 415 //Intentionally printing dates in their own timezone, not Oozie timezone 416 LOG.info("materialize actions for tz=" + appTz.getDisplayName() + ",\n start=" + start.getTime() + ", end=" 417 + end.getTime() + ",\n timeUnit " + freqTU.getCalendarUnit() + ",\n frequency :" + frequency + ":" 418 + freqTU + ",\n lastActionNumber " + lastActionNumber); 419 // Keep the actual start time 420 Calendar origStart = Calendar.getInstance(appTz); 421 origStart.setTime(coordJob.getStartTimestamp()); 422 // Move to the End of duration, if needed. 423 DateUtils.moveToEnd(origStart, endOfFlag); 424 425 StringBuilder actionStrings = new StringBuilder(); 426 Date jobPauseTime = coordJob.getPauseTime(); 427 Calendar pause = null; 428 if (jobPauseTime != null) { 429 pause = Calendar.getInstance(appTz); 430 pause.setTime(DateUtils.convertDateToTimestamp(jobPauseTime)); 431 } 432 433 String action = null; 434 int numWaitingActions = dryrun ? 0 : jpaService.execute(new CoordActionsActiveCountJPAExecutor(coordJob.getId())); 435 int maxActionToBeCreated = coordJob.getMatThrottling() - numWaitingActions; 436 // If LAST_ONLY and all materialization is in the past, ignore maxActionsToBeCreated 437 boolean ignoreMaxActions = 438 (coordJob.getExecutionOrder().equals(CoordinatorJob.Execution.LAST_ONLY) || 439 coordJob.getExecutionOrder().equals(CoordinatorJob.Execution.NONE)) 440 && endMatdTime.before(new Date()); 441 LOG.debug("Coordinator job :" + coordJob.getId() + ", maxActionToBeCreated :" + maxActionToBeCreated 442 + ", Mat_Throttle :" + coordJob.getMatThrottling() + ", numWaitingActions :" + numWaitingActions); 443 444 boolean isCronFrequency = false; 445 446 Calendar effStart = (Calendar) start.clone(); 447 try { 448 int intFrequency = Integer.parseInt(coordJob.getFrequency()); 449 effStart = (Calendar) origStart.clone(); 450 effStart.add(freqTU.getCalendarUnit(), lastActionNumber * intFrequency); 451 } 452 catch (NumberFormatException e) { 453 isCronFrequency = true; 454 } 455 456 boolean firstMater = true; 457 while (effStart.compareTo(end) < 0 && (ignoreMaxActions || maxActionToBeCreated-- > 0)) { 458 if (pause != null && effStart.compareTo(pause) >= 0) { 459 break; 460 } 461 462 Date nextTime = effStart.getTime(); 463 464 if (isCronFrequency) { 465 if (effStart.getTime().compareTo(startMatdTime) == 0 && firstMater) { 466 effStart.add(Calendar.MINUTE, -1); 467 firstMater = false; 468 } 469 470 nextTime = CoordCommandUtils.getNextValidActionTimeForCronFrequency(effStart.getTime(), coordJob); 471 effStart.setTime(nextTime); 472 } 473 474 if (effStart.compareTo(end) < 0) { 475 476 if (pause != null && effStart.compareTo(pause) >= 0) { 477 break; 478 } 479 CoordinatorActionBean actionBean = new CoordinatorActionBean(); 480 lastActionNumber++; 481 482 int timeout = coordJob.getTimeout(); 483 LOG.debug("Materializing action for time=" + DateUtils.formatDateOozieTZ(effStart.getTime()) 484 + ", lastactionnumber=" + lastActionNumber + " timeout=" + timeout + " minutes"); 485 Date actualTime = new Date(); 486 action = CoordCommandUtils.materializeOneInstance(jobId, dryrun, (Element) eJob.clone(), 487 nextTime, actualTime, lastActionNumber, jobConf, actionBean); 488 actionBean.setTimeOut(timeout); 489 490 if (!dryrun) { 491 storeToDB(actionBean, action); // Storing to table 492 493 } 494 else { 495 actionStrings.append("action for new instance"); 496 actionStrings.append(action); 497 } 498 } 499 else { 500 break; 501 } 502 503 if (!isCronFrequency) { 504 effStart = (Calendar) origStart.clone(); 505 effStart.add(freqTU.getCalendarUnit(), lastActionNumber * Integer.parseInt(coordJob.getFrequency())); 506 } 507 } 508 509 if (isCronFrequency) { 510 if (effStart.compareTo(end) < 0 && !(ignoreMaxActions || maxActionToBeCreated-- > 0)) { 511 //Since we exceed the throttle, we need to move the nextMadtime forward 512 //to avoid creating duplicate actions 513 if (!firstMater) { 514 effStart.setTime(CoordCommandUtils.getNextValidActionTimeForCronFrequency(effStart.getTime(), coordJob)); 515 } 516 } 517 } 518 519 endMatdTime = effStart.getTime(); 520 521 if (!dryrun) { 522 return action; 523 } 524 else { 525 return actionStrings.toString(); 526 } 527 } 528 529 private void storeToDB(CoordinatorActionBean actionBean, String actionXml) throws Exception { 530 LOG.debug("In storeToDB() coord action id = " + actionBean.getId() + ", size of actionXml = " 531 + actionXml.length()); 532 actionBean.setActionXml(actionXml); 533 534 insertList.add(actionBean); 535 writeActionSlaRegistration(actionXml, actionBean); 536 } 537 538 private void writeActionSlaRegistration(String actionXml, CoordinatorActionBean actionBean) throws Exception { 539 Element eAction = XmlUtils.parseXml(actionXml); 540 Element eSla = eAction.getChild("action", eAction.getNamespace()).getChild("info", eAction.getNamespace("sla")); 541 SLAEventBean slaEvent = SLADbOperations.createSlaRegistrationEvent(eSla, actionBean.getId(), SlaAppType.COORDINATOR_ACTION, coordJob 542 .getUser(), coordJob.getGroup(), LOG); 543 if(slaEvent != null) { 544 insertList.add(slaEvent); 545 } 546 // inserting into new table also 547 SLAOperations.createSlaRegistrationEvent(eSla, actionBean.getId(), actionBean.getJobId(), 548 AppType.COORDINATOR_ACTION, coordJob.getUser(), coordJob.getAppName(), LOG, false); 549 } 550 551 private void updateJobMaterializeInfo(CoordinatorJobBean job) throws CommandException { 552 job.setLastActionTime(endMatdTime); 553 job.setLastActionNumber(lastActionNumber); 554 // if the job endtime == action endtime, we don't need to materialize this job anymore 555 Date jobEndTime = job.getEndTime(); 556 557 558 if (job.getStatus() == CoordinatorJob.Status.PREP){ 559 LOG.info("[" + job.getId() + "]: Update status from " + job.getStatus() + " to RUNNING"); 560 job.setStatus(Job.Status.RUNNING); 561 } 562 job.setPending(); 563 564 if (jobEndTime.compareTo(endMatdTime) <= 0) { 565 LOG.info("[" + job.getId() + "]: all actions have been materialized, set pending to true"); 566 // set doneMaterialization to true when materialization is done 567 job.setDoneMaterialization(); 568 } 569 job.setStatus(StatusUtils.getStatus(job)); 570 LOG.info("Coord Job status updated to = " + job.getStatus()); 571 job.setNextMaterializedTime(endMatdTime); 572 } 573 574 /* (non-Javadoc) 575 * @see org.apache.oozie.command.XCommand#getKey() 576 */ 577 @Override 578 public String getKey() { 579 return getName() + "_" + jobId; 580 } 581 582 /* (non-Javadoc) 583 * @see org.apache.oozie.command.TransitionXCommand#notifyParent() 584 */ 585 @Override 586 public void notifyParent() throws CommandException { 587 // update bundle action only when status changes in coord job 588 if (this.coordJob.getBundleId() != null) { 589 if (!prevStatus.equals(coordJob.getStatus())) { 590 BundleStatusUpdateXCommand bundleStatusUpdate = new BundleStatusUpdateXCommand(coordJob, prevStatus); 591 bundleStatusUpdate.call(); 592 } 593 } 594 } 595}