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.sla; 020 021import java.sql.Timestamp; 022import java.util.ArrayList; 023import java.util.Collections; 024import java.util.Date; 025import java.util.HashSet; 026import java.util.Iterator; 027import java.util.List; 028import java.util.Map; 029import java.util.Set; 030import java.util.concurrent.ConcurrentHashMap; 031 032import org.apache.hadoop.conf.Configuration; 033import org.apache.oozie.AppType; 034import org.apache.oozie.CoordinatorActionBean; 035import org.apache.oozie.CoordinatorJobBean; 036import org.apache.oozie.ErrorCode; 037import org.apache.oozie.WorkflowActionBean; 038import org.apache.oozie.WorkflowJobBean; 039import org.apache.oozie.XException; 040import org.apache.oozie.client.CoordinatorAction; 041import org.apache.oozie.client.WorkflowAction; 042import org.apache.oozie.client.WorkflowJob; 043import org.apache.oozie.client.event.JobEvent; 044import org.apache.oozie.client.event.SLAEvent.EventStatus; 045import org.apache.oozie.client.event.SLAEvent.SLAStatus; 046import org.apache.oozie.client.rest.JsonBean; 047import org.apache.oozie.executor.jpa.BatchQueryExecutor; 048import org.apache.oozie.executor.jpa.CoordActionGetForSLAJPAExecutor; 049import org.apache.oozie.executor.jpa.CoordActionQueryExecutor; 050import org.apache.oozie.executor.jpa.CoordActionQueryExecutor.CoordActionQuery; 051import org.apache.oozie.executor.jpa.CoordJobQueryExecutor; 052import org.apache.oozie.executor.jpa.CoordJobQueryExecutor.CoordJobQuery; 053import org.apache.oozie.executor.jpa.JPAExecutorException; 054import org.apache.oozie.executor.jpa.SLARegistrationQueryExecutor; 055import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor; 056import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor; 057import org.apache.oozie.executor.jpa.SLARegistrationQueryExecutor.SLARegQuery; 058import org.apache.oozie.executor.jpa.SLASummaryQueryExecutor; 059import org.apache.oozie.executor.jpa.WorkflowActionGetForSLAJPAExecutor; 060import org.apache.oozie.executor.jpa.WorkflowJobGetForSLAJPAExecutor; 061import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor.WorkflowActionQuery; 062import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor.WorkflowJobQuery; 063import org.apache.oozie.executor.jpa.sla.SLASummaryGetRecordsOnRestartJPAExecutor; 064import org.apache.oozie.executor.jpa.SLASummaryQueryExecutor.SLASummaryQuery; 065import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry; 066import org.apache.oozie.lock.LockToken; 067import org.apache.oozie.service.ConfigurationService; 068import org.apache.oozie.service.EventHandlerService; 069import org.apache.oozie.service.InstrumentationService; 070import org.apache.oozie.service.JPAService; 071import org.apache.oozie.service.JobsConcurrencyService; 072import org.apache.oozie.service.MemoryLocksService; 073import org.apache.oozie.service.SchedulerService; 074import org.apache.oozie.service.ServiceException; 075import org.apache.oozie.service.Services; 076import org.apache.oozie.sla.service.SLAService; 077import org.apache.oozie.util.DateUtils; 078import org.apache.oozie.util.Instrumentation; 079import org.apache.oozie.util.LogUtils; 080import org.apache.oozie.util.XLog; 081 082import com.google.common.annotations.VisibleForTesting; 083 084 085/** 086 * Implementation class for SLACalculator that calculates SLA related to 087 * start/end/duration of jobs using a memory-based map 088 */ 089public class SLACalculatorMemory implements SLACalculator { 090 091 private static XLog LOG = XLog.getLog(SLACalculatorMemory.class); 092 // TODO optimization priority based insertion/processing/bumping up-down 093 private Map<String, SLACalcStatus> slaMap; 094 protected Set<String> historySet; 095 private static int capacity; 096 private static JPAService jpaService; 097 protected EventHandlerService eventHandler; 098 private static int modifiedAfter; 099 private static long jobEventLatency; 100 private Instrumentation instrumentation; 101 public static final String INSTRUMENTATION_GROUP = "sla-calculator"; 102 public static final String SLA_MAP = "sla-map"; 103 private int maxRetryCount; 104 105 @Override 106 public void init(Configuration conf) throws ServiceException { 107 capacity = ConfigurationService.getInt(conf, SLAService.CONF_CAPACITY); 108 jobEventLatency = ConfigurationService.getInt(conf, SLAService.CONF_JOB_EVENT_LATENCY); 109 maxRetryCount = ConfigurationService.getInt(conf, SLAService.CONF_MAXIMUM_RETRY_COUNT); 110 slaMap = new ConcurrentHashMap<String, SLACalcStatus>(); 111 historySet = Collections.synchronizedSet(new HashSet<String>()); 112 jpaService = Services.get().get(JPAService.class); 113 eventHandler = Services.get().get(EventHandlerService.class); 114 instrumentation = Services.get().get(InstrumentationService.class).get(); 115 // load events modified after 116 modifiedAfter = conf.getInt(SLAService.CONF_EVENTS_MODIFIED_AFTER, 7); 117 loadOnRestart(); 118 Runnable purgeThread = new HistoryPurgeWorker(); 119 // schedule runnable by default 1 day 120 Services.get() 121 .get(SchedulerService.class) 122 .schedule(purgeThread, 86400, Services.get().getConf().getInt(SLAService.CONF_SLA_HISTORY_PURGE_INTERVAL, 86400), 123 SchedulerService.Unit.SEC); 124 } 125 126 public class HistoryPurgeWorker implements Runnable { 127 128 public HistoryPurgeWorker() { 129 } 130 131 @Override 132 public void run() { 133 if (Thread.currentThread().isInterrupted()) { 134 return; 135 } 136 Iterator<String> jobItr = historySet.iterator(); 137 while (jobItr.hasNext()) { 138 String jobId = jobItr.next(); 139 140 if (jobId.endsWith("-W")) { 141 WorkflowJobBean wfJob = null; 142 try { 143 wfJob = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW_STATUS, jobId); 144 } 145 catch (JPAExecutorException e) { 146 if (e.getErrorCode().equals(ErrorCode.E0604)) { 147 jobItr.remove(); 148 } 149 else { 150 LOG.info("Failed to fetch the workflow job: " + jobId, e); 151 } 152 } 153 if (wfJob != null && wfJob.inTerminalState()) { 154 try { 155 updateSLASummary(wfJob.getId(), wfJob.getStartTime(), wfJob.getEndTime()); 156 jobItr.remove(); 157 } 158 catch (JPAExecutorException e) { 159 LOG.info("Failed to update SLASummaryBean when purging history set entry for " + jobId, e); 160 } 161 162 } 163 } 164 else if (jobId.contains("-W@")) { 165 WorkflowActionBean wfAction = null; 166 try { 167 wfAction = WorkflowActionQueryExecutor.getInstance().get( 168 WorkflowActionQuery.GET_ACTION_COMPLETED, jobId); 169 } 170 catch (JPAExecutorException e) { 171 if (e.getErrorCode().equals(ErrorCode.E0605)) { 172 jobItr.remove(); 173 } 174 else { 175 LOG.info("Failed to fetch the workflow action: " + jobId, e); 176 } 177 } 178 if (wfAction != null && (wfAction.isComplete() || wfAction.isTerminalWithFailure())) { 179 try { 180 updateSLASummary(wfAction.getId(), wfAction.getStartTime(), wfAction.getEndTime()); 181 jobItr.remove(); 182 } 183 catch (JPAExecutorException e) { 184 LOG.info("Failed to update SLASummaryBean when purging history set entry for " + jobId, e); 185 } 186 } 187 } 188 else if (jobId.contains("-C@")) { 189 CoordinatorActionBean cAction = null; 190 try { 191 cAction = CoordActionQueryExecutor.getInstance().get(CoordActionQuery.GET_COORD_ACTION, jobId); 192 } 193 catch (JPAExecutorException e) { 194 if (e.getErrorCode().equals(ErrorCode.E0605)) { 195 jobItr.remove(); 196 } 197 else { 198 LOG.info("Failed to fetch the coord action: " + jobId, e); 199 } 200 } 201 if (cAction != null && cAction.isTerminalStatus()) { 202 try { 203 updateSLASummaryForCoordAction(cAction); 204 jobItr.remove(); 205 } 206 catch (JPAExecutorException e) { 207 XLog.getLog(SLACalculatorMemory.class).info( 208 "Failed to update SLASummaryBean when purging history set entry for " + jobId, e); 209 } 210 211 } 212 } 213 else if (jobId.endsWith("-C")) { 214 CoordinatorJobBean cJob = null; 215 try { 216 cJob = CoordJobQueryExecutor.getInstance().get(CoordJobQuery.GET_COORD_JOB_STATUS_PARENTID, 217 jobId); 218 } 219 catch (JPAExecutorException e) { 220 if (e.getErrorCode().equals(ErrorCode.E0604)) { 221 jobItr.remove(); 222 } 223 else { 224 LOG.info("Failed to fetch the coord job: " + jobId, e); 225 } 226 } 227 if (cJob != null && cJob.isTerminalStatus()) { 228 try { 229 updateSLASummary(cJob.getId(), cJob.getStartTime(), cJob.getEndTime()); 230 jobItr.remove(); 231 } 232 catch (JPAExecutorException e) { 233 LOG.info("Failed to update SLASummaryBean when purging history set entry for " + jobId, e); 234 } 235 236 } 237 } 238 } 239 } 240 241 private void updateSLASummary(String id, Date startTime, Date endTime) throws JPAExecutorException { 242 SLASummaryBean sla = SLASummaryQueryExecutor.getInstance().get(SLASummaryQuery.GET_SLA_SUMMARY, id); 243 if (sla != null) { 244 sla.setActualStart(startTime); 245 sla.setActualEnd(endTime); 246 if (startTime != null && endTime != null) { 247 sla.setActualDuration(endTime.getTime() - startTime.getTime()); 248 } 249 sla.setLastModifiedTime(new Date()); 250 sla.setEventProcessed(8); 251 SLASummaryQueryExecutor.getInstance().executeUpdate( 252 SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_ACTUAL_TIMES, sla); 253 } 254 } 255 256 private void updateSLASummaryForCoordAction(CoordinatorActionBean bean) throws JPAExecutorException { 257 String wrkflowId = bean.getExternalId(); 258 if (wrkflowId != null) { 259 WorkflowJobBean wrkflow = WorkflowJobQueryExecutor.getInstance().get( 260 WorkflowJobQuery.GET_WORKFLOW_START_END_TIME, wrkflowId); 261 if (wrkflow != null) { 262 updateSLASummary(bean.getId(), wrkflow.getStartTime(), wrkflow.getEndTime()); 263 } 264 } 265 } 266 } 267 268 private void loadOnRestart() { 269 boolean isJobModified = false; 270 try { 271 long slaPendingCount = 0; 272 long statusPendingCount = 0; 273 List<SLASummaryBean> summaryBeans = jpaService.execute(new SLASummaryGetRecordsOnRestartJPAExecutor( 274 modifiedAfter)); 275 for (SLASummaryBean summaryBean : summaryBeans) { 276 String jobId = summaryBean.getId(); 277 LockToken lock = null; 278 try { 279 switch (summaryBean.getAppType()) { 280 case COORDINATOR_ACTION: 281 isJobModified = processSummaryBeanForCoordAction(summaryBean, jobId); 282 break; 283 case WORKFLOW_ACTION: 284 isJobModified = processSummaryBeanForWorkflowAction(summaryBean, jobId); 285 break; 286 case WORKFLOW_JOB: 287 isJobModified = processSummaryBeanForWorkflowJob(summaryBean, jobId); 288 break; 289 default: 290 break; 291 } 292 } 293 catch (final JPAExecutorException | IllegalArgumentException e) { 294 LOG.warn("Failed to process SLA summary bean " + jobId, e); 295 } 296 if (isJobModified) { 297 try { 298 boolean update = true; 299 if (Services.get().get(JobsConcurrencyService.class).isHighlyAvailableMode()) { 300 lock = Services 301 .get() 302 .get(MemoryLocksService.class) 303 .getWriteLock( 304 SLACalcStatus.SLA_ENTITYKEY_PREFIX + jobId, 305 Services.get().getConf() 306 .getLong(SLAService.CONF_SLA_CALC_LOCK_TIMEOUT, 5 * 1000)); 307 if (lock == null) { 308 update = false; 309 } 310 } 311 if (update) { 312 summaryBean.setLastModifiedTime(new Date()); 313 SLASummaryQueryExecutor.getInstance().executeUpdate( 314 SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_STATUS_ACTUAL_TIMES, summaryBean); 315 } 316 } 317 catch (Exception e) { 318 LOG.warn("Failed to load records for " + jobId, e); 319 } 320 finally { 321 if (lock != null) { 322 lock.release(); 323 lock = null; 324 } 325 } 326 } 327 try { 328 if (summaryBean.getEventProcessed() == 7) { 329 historySet.add(jobId); 330 statusPendingCount++; 331 } 332 else if (summaryBean.getEventProcessed() <= 7) { 333 SLARegistrationBean slaRegBean = SLARegistrationQueryExecutor.getInstance().get( 334 SLARegQuery.GET_SLA_REG_ON_RESTART, jobId); 335 SLACalcStatus slaCalcStatus = new SLACalcStatus(summaryBean, slaRegBean); 336 putAndIncrement(jobId, slaCalcStatus); 337 slaPendingCount++; 338 } 339 } 340 catch (Exception e) { 341 LOG.warn("Failed to fetch/update records for " + jobId, e); 342 } 343 344 } 345 LOG.info("Loaded SLASummary pendingSLA=" + slaPendingCount + ", pendingStatusUpdate=" + statusPendingCount); 346 347 } 348 catch (Exception e) { 349 LOG.warn("Failed to retrieve SLASummary records on restart", e); 350 } 351 } 352 353 private boolean processSummaryBeanForCoordAction(SLASummaryBean summaryBean, String jobId) 354 throws JPAExecutorException { 355 boolean isJobModified = false; 356 CoordinatorActionBean coordAction = null; 357 coordAction = jpaService.execute(new CoordActionGetForSLAJPAExecutor(jobId)); 358 if (!coordAction.getStatusStr().equals(summaryBean.getJobStatus())) { 359 LOG.trace("Coordinator action status is " + coordAction.getStatusStr() + " and summary bean status is " 360 + summaryBean.getJobStatus()); 361 isJobModified = true; 362 summaryBean.setJobStatus(coordAction.getStatusStr()); 363 if (coordAction.isTerminalStatus()) { 364 WorkflowJobBean wfJob = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(coordAction 365 .getExternalId())); 366 setEndForSLASummaryBean(summaryBean, wfJob.getStartTime(), coordAction.getLastModifiedTime(), 367 coordAction.getStatusStr()); 368 } 369 else if (coordAction.getStatus() != CoordinatorAction.Status.WAITING) { 370 WorkflowJobBean wfJob = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(coordAction 371 .getExternalId())); 372 setStartForSLASummaryBean(summaryBean, summaryBean.getEventProcessed(), wfJob.getStartTime()); 373 } 374 } 375 return isJobModified; 376 } 377 378 private boolean processSummaryBeanForWorkflowAction(SLASummaryBean summaryBean, String jobId) 379 throws JPAExecutorException { 380 boolean isJobModified = false; 381 WorkflowActionBean wfAction = null; 382 wfAction = jpaService.execute(new WorkflowActionGetForSLAJPAExecutor(jobId)); 383 if (!wfAction.getStatusStr().equals(summaryBean.getJobStatus())) { 384 LOG.trace("Workflow action status is " + wfAction.getStatusStr() + "and summary bean status is " 385 + summaryBean.getJobStatus()); 386 isJobModified = true; 387 summaryBean.setJobStatus(wfAction.getStatusStr()); 388 if (wfAction.inTerminalState()) { 389 setEndForSLASummaryBean(summaryBean, wfAction.getStartTime(), wfAction.getEndTime(), wfAction.getStatusStr()); 390 } 391 else if (wfAction.getStatus() != WorkflowAction.Status.PREP) { 392 setStartForSLASummaryBean(summaryBean, summaryBean.getEventProcessed(), wfAction.getStartTime()); 393 } 394 } 395 return isJobModified; 396 } 397 398 private boolean processSummaryBeanForWorkflowJob(SLASummaryBean summaryBean, String jobId) 399 throws JPAExecutorException { 400 boolean isJobModified = false; 401 WorkflowJobBean wfJob = null; 402 wfJob = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(jobId)); 403 if (!wfJob.getStatusStr().equals(summaryBean.getJobStatus())) { 404 LOG.trace("Workflow job status is " + wfJob.getStatusStr() + "and summary bean status is " 405 + summaryBean.getJobStatus()); 406 isJobModified = true; 407 summaryBean.setJobStatus(wfJob.getStatusStr()); 408 if (wfJob.inTerminalState()) { 409 setEndForSLASummaryBean(summaryBean, wfJob.getStartTime(), wfJob.getEndTime(), wfJob.getStatusStr()); 410 } 411 else if (wfJob.getStatus() != WorkflowJob.Status.PREP) { 412 setStartForSLASummaryBean(summaryBean, summaryBean.getEventProcessed(), wfJob.getStartTime()); 413 } 414 } 415 return isJobModified; 416 } 417 418 private void setEndForSLASummaryBean(SLASummaryBean summaryBean, Date startTime, Date endTime, String status) { 419 byte eventProc = summaryBean.getEventProcessed(); 420 summaryBean.setEventProcessed(8); 421 summaryBean.setActualStart(startTime); 422 summaryBean.setActualEnd(endTime); 423 long actualDuration = endTime.getTime() - startTime.getTime(); 424 summaryBean.setActualDuration(actualDuration); 425 if (eventProc < 4) { 426 if (status.equals(WorkflowJob.Status.SUCCEEDED.name()) || status.equals(WorkflowAction.Status.OK.name()) 427 || status.equals(CoordinatorAction.Status.SUCCEEDED.name())) { 428 if (endTime.getTime() <= summaryBean.getExpectedEnd().getTime()) { 429 summaryBean.setSLAStatus(SLAStatus.MET); 430 } 431 else { 432 summaryBean.setSLAStatus(SLAStatus.MISS); 433 } 434 } 435 else { 436 summaryBean.setSLAStatus(SLAStatus.MISS); 437 } 438 } 439 440 } 441 442 private void setStartForSLASummaryBean(SLASummaryBean summaryBean, byte eventProc, Date startTime) { 443 if (((eventProc & 1) == 0)) { 444 eventProc += 1; 445 summaryBean.setEventProcessed(eventProc); 446 } 447 if (summaryBean.getSLAStatus().equals(SLAStatus.NOT_STARTED)) { 448 summaryBean.setSLAStatus(SLAStatus.IN_PROCESS); 449 } 450 summaryBean.setActualStart(startTime); 451 } 452 453 @Override 454 public int size() { 455 return slaMap.size(); 456 } 457 458 @Override 459 public SLACalcStatus get(String jobId) throws JPAExecutorException { 460 SLACalcStatus memObj; 461 memObj = slaMap.get(jobId); 462 if (memObj == null && historySet.contains(jobId)) { 463 memObj = new SLACalcStatus(SLASummaryQueryExecutor.getInstance().get(SLASummaryQuery.GET_SLA_SUMMARY, jobId), 464 SLARegistrationQueryExecutor.getInstance().get(SLARegQuery.GET_SLA_REG_ON_RESTART, jobId)); 465 } 466 return memObj; 467 } 468 469 @Override 470 public Iterator<String> iterator() { 471 return slaMap.keySet().iterator(); 472 } 473 474 @Override 475 public boolean isEmpty() { 476 return slaMap.isEmpty(); 477 } 478 479 @Override 480 public void clear() { 481 final int originalSize = slaMap.size(); 482 LOG.trace("Clearing SLA map and SLA history set. [slaMap.size={};slaHistorySet.size={1}]", 483 slaMap.size(), 484 historySet.size()); 485 slaMap.clear(); 486 historySet.clear(); 487 instrumentation.decr(INSTRUMENTATION_GROUP, SLA_MAP, originalSize); 488 } 489 490 /** 491 * Invoked via periodic run, update the SLA for registered jobs. 492 * <p> 493 * Track the number of times the {@link SLACalcStatus} entry has not been processed successfully, and when a preconfigured 494 * {code oozie.sla.service.SLAService.maximum.retry.count} is reached, remove any {@link SLACalculatorMemory#slaMap} entries 495 * that are causing {@code JPAExecutorException}s of certain {@link ErrorCode}s. 496 * @param jobId the workflow or coordinator job or action ID the SLA is tracked against 497 */ 498 void updateJobSla(String jobId) throws Exception { 499 SLACalcStatus slaCalc = slaMap.get(jobId); 500 synchronized (slaCalc) { 501 boolean change = false; 502 Object eventProcObj = null; 503 try { 504 // get eventProcessed on DB for validation in HA 505 eventProcObj = ((SLASummaryQueryExecutor) SLASummaryQueryExecutor.getInstance()).getSingleValue( 506 SLASummaryQuery.GET_SLA_SUMMARY_EVENTPROCESSED, jobId); 507 resetRetryCount(jobId); 508 } catch (final JPAExecutorException e) { 509 if (e.getErrorCode().equals(ErrorCode.E0603) 510 || e.getErrorCode().equals(ErrorCode.E0604) 511 || e.getErrorCode().equals(ErrorCode.E0605)) { 512 LOG.debug("job [{0}] is not in DB, removing from Memory", jobId); 513 incrementRetryCountAndRemove(jobId); 514 return; 515 } 516 throw e; 517 } 518 byte eventProc = ((Byte) eventProcObj).byteValue(); 519 if (eventProc >= 7) { 520 if (eventProc == 7) { 521 historySet.add(jobId); 522 } 523 removeAndDecrement(jobId); 524 LOG.trace("Removed Job [{0}] from map as SLA processed", jobId); 525 } 526 else { 527 slaCalc.setEventProcessed(eventProc); 528 SLARegistrationBean reg = slaCalc.getSLARegistrationBean(); 529 // calculation w.r.t current time and status 530 if ((eventProc & 1) == 0) { // first bit (start-processed) unset 531 if (reg.getExpectedStart() != null) { 532 if (reg.getExpectedStart().getTime() + jobEventLatency < System.currentTimeMillis()) { 533 confirmWithDB(slaCalc); 534 eventProc = slaCalc.getEventProcessed(); 535 if (eventProc != 8 && (eventProc & 1) == 0) { 536 // Some DB exception 537 slaCalc.setEventStatus(EventStatus.START_MISS); 538 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 539 eventProc++; 540 } 541 change = true; 542 } 543 } 544 else { 545 eventProc++; // disable further processing for optional start sla condition 546 change = true; 547 } 548 } 549 // check if second bit (duration-processed) is unset 550 if (eventProc != 8 && ((eventProc >> 1) & 1) == 0) { 551 if (reg.getExpectedDuration() == -1) { 552 eventProc += 2; 553 change = true; 554 } 555 else if (slaCalc.getActualStart() != null) { 556 if ((reg.getExpectedDuration() + jobEventLatency) < (System.currentTimeMillis() - slaCalc 557 .getActualStart().getTime())) { 558 slaCalc.setEventProcessed(eventProc); 559 confirmWithDB(slaCalc); 560 eventProc = slaCalc.getEventProcessed(); 561 if (eventProc != 8 && ((eventProc >> 1) & 1) == 0) { 562 // Some DB exception 563 slaCalc.setEventStatus(EventStatus.DURATION_MISS); 564 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 565 eventProc += 2; 566 } 567 change = true; 568 } 569 } 570 } 571 if (eventProc < 4) { 572 if (reg.getExpectedEnd().getTime() + jobEventLatency < System.currentTimeMillis()) { 573 slaCalc.setEventProcessed(eventProc); 574 confirmWithDB(slaCalc); 575 eventProc = slaCalc.getEventProcessed(); 576 change = true; 577 } 578 } 579 if (change) { 580 try { 581 boolean locked = true; 582 slaCalc.acquireLock(); 583 locked = slaCalc.isLocked(); 584 if (locked) { 585 // no more processing, no transfer to history set 586 if (slaCalc.getEventProcessed() >= 8) { 587 eventProc = 8; 588 // Should not be > 8. But to handle any corner cases 589 slaCalc.setEventProcessed(8); 590 removeAndDecrement(jobId); 591 } 592 else { 593 slaCalc.setEventProcessed(eventProc); 594 } 595 SLASummaryBean slaSummaryBean = new SLASummaryBean(); 596 slaSummaryBean.setId(slaCalc.getId()); 597 slaSummaryBean.setEventProcessed(eventProc); 598 slaSummaryBean.setSLAStatus(slaCalc.getSLAStatus()); 599 slaSummaryBean.setEventStatus(slaCalc.getEventStatus()); 600 slaSummaryBean.setActualEnd(slaCalc.getActualEnd()); 601 slaSummaryBean.setActualStart(slaCalc.getActualStart()); 602 slaSummaryBean.setActualDuration(slaCalc.getActualDuration()); 603 slaSummaryBean.setJobStatus(slaCalc.getJobStatus()); 604 slaSummaryBean.setLastModifiedTime(new Date()); 605 SLASummaryQueryExecutor.getInstance().executeUpdate( 606 SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_STATUS_ACTUAL_TIMES, slaSummaryBean); 607 if (eventProc == 7) { 608 historySet.add(jobId); 609 removeAndDecrement(jobId); 610 LOG.trace("Removed Job [{0}] from map after End-processed", jobId); 611 } 612 } 613 } 614 catch (InterruptedException e) { 615 incrementRetryCountAndRemove(jobId); 616 throw new XException(ErrorCode.E0606, slaCalc.getId(), slaCalc.getLockTimeOut()); 617 } 618 finally { 619 slaCalc.releaseLock(); 620 } 621 } 622 } 623 } 624 } 625 626 /** 627 * Periodically run by the SLAService worker threads to update SLA status by 628 * iterating through all the jobs in the map 629 */ 630 @Override 631 public void updateAllSlaStatus() { 632 LOG.info("Running periodic SLA check"); 633 Iterator<String> iterator = slaMap.keySet().iterator(); 634 while (iterator.hasNext()) { 635 String jobId = iterator.next(); 636 try { 637 LOG.trace("Processing SLA for jobid={0}", jobId); 638 updateJobSla(jobId); 639 } 640 catch (Exception e) { 641 setLogPrefix(jobId); 642 LOG.error("Exception in SLA processing for job [{0}]", jobId, e); 643 LogUtils.clearLogPrefix(); 644 } 645 } 646 } 647 648 /** 649 * Register a new job into the map for SLA tracking 650 */ 651 @Override 652 public boolean addRegistration(String jobId, SLARegistrationBean reg) throws JPAExecutorException { 653 try { 654 if (slaMap.size() < capacity) { 655 SLACalcStatus slaCalc = new SLACalcStatus(reg); 656 slaCalc.setSLAStatus(SLAStatus.NOT_STARTED); 657 slaCalc.setJobStatus(getJobStatus(reg.getAppType())); 658 List<JsonBean> insertList = new ArrayList<JsonBean>(); 659 final SLASummaryBean summaryBean = new SLASummaryBean(slaCalc); 660 final Timestamp currentTime = DateUtils.convertDateToTimestamp(new Date()); 661 reg.setCreatedTimestamp(currentTime); 662 summaryBean.setCreatedTimestamp(currentTime); 663 insertList.add(reg); 664 insertList.add(summaryBean); 665 BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, null, null); 666 putAndIncrement(jobId, slaCalc); 667 LOG.trace("SLA Registration Event - Job:" + jobId); 668 return true; 669 } 670 else { 671 setLogPrefix(reg.getId()); 672 LOG.error( 673 "SLACalculator memory capacity reached. Cannot add or update new SLA Registration entry for job [{0}]", 674 reg.getId()); 675 LogUtils.clearLogPrefix(); 676 } 677 } 678 catch (JPAExecutorException jpa) { 679 throw jpa; 680 } 681 return false; 682 } 683 684 private String getJobStatus(AppType appType) { 685 String status = null; 686 switch (appType) { 687 case COORDINATOR_ACTION: 688 status = CoordinatorAction.Status.WAITING.name(); 689 break; 690 case WORKFLOW_ACTION: 691 status = WorkflowAction.Status.PREP.name(); 692 break; 693 case WORKFLOW_JOB: 694 status = WorkflowJob.Status.PREP.name(); 695 break; 696 default: 697 break; 698 } 699 return status; 700 } 701 702 /** 703 * Update job into the map for SLA tracking 704 */ 705 @Override 706 public boolean updateRegistration(String jobId, SLARegistrationBean reg) throws JPAExecutorException { 707 try { 708 if (slaMap.size() < capacity) { 709 SLACalcStatus slaCalc = new SLACalcStatus(reg); 710 slaCalc.setSLAStatus(SLAStatus.NOT_STARTED); 711 slaCalc.setJobStatus(getJobStatus(reg.getAppType())); 712 @SuppressWarnings("rawtypes") 713 List<UpdateEntry> updateList = new ArrayList<UpdateEntry>(); 714 updateList.add(new UpdateEntry<SLARegQuery>(SLARegQuery.UPDATE_SLA_REG_ALL, reg)); 715 updateList.add(new UpdateEntry<SLASummaryQuery>(SLASummaryQuery.UPDATE_SLA_SUMMARY_ALL, 716 new SLASummaryBean(slaCalc))); 717 BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(null, updateList, null); 718 putAndIncrement(jobId, slaCalc); 719 LOG.trace("SLA Registration Event - Job:" + jobId); 720 return true; 721 } 722 else { 723 setLogPrefix(reg.getId()); 724 LOG.error( 725 "SLACalculator memory capacity reached. Cannot add or update new SLA Registration entry for job [{0}]", 726 reg.getId()); 727 LogUtils.clearLogPrefix(); 728 } 729 } 730 catch (JPAExecutorException jpa) { 731 throw jpa; 732 } 733 return false; 734 } 735 736 /** 737 * Remove job from being tracked in map 738 */ 739 @Override 740 public void removeRegistration(String jobId) { 741 if (!removeAndDecrement(jobId)) { 742 historySet.remove(jobId); 743 } 744 } 745 746 /** 747 * Triggered after receiving Job status change event, update SLA status 748 * accordingly 749 */ 750 @Override 751 public boolean addJobStatus(String jobId, String jobStatus, JobEvent.EventStatus jobEventStatus, Date startTime, 752 Date endTime) throws JPAExecutorException, ServiceException { 753 SLACalcStatus slaCalc = slaMap.get(jobId); 754 SLASummaryBean slaInfo = null; 755 boolean hasSla = false; 756 if (slaCalc == null) { 757 if (historySet.contains(jobId)) { 758 slaInfo = SLASummaryQueryExecutor.getInstance().get(SLASummaryQuery.GET_SLA_SUMMARY, jobId); 759 if (slaInfo == null) { 760 throw new JPAExecutorException(ErrorCode.E0604, jobId); 761 } 762 slaInfo.setJobStatus(jobStatus); 763 slaInfo.setActualStart(startTime); 764 slaInfo.setActualEnd(endTime); 765 if (endTime != null) { 766 slaInfo.setActualDuration(endTime.getTime() - startTime.getTime()); 767 } 768 slaInfo.setEventProcessed(8); 769 historySet.remove(jobId); 770 slaInfo.setLastModifiedTime(new Date()); 771 SLASummaryQueryExecutor.getInstance().executeUpdate( 772 SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_STATUS_ACTUAL_TIMES, slaInfo); 773 hasSla = true; 774 } 775 else if (Services.get().get(JobsConcurrencyService.class).isHighlyAvailableMode()) { 776 // jobid might not exist in slaMap in HA Setting 777 SLARegistrationBean slaRegBean = SLARegistrationQueryExecutor.getInstance().get( 778 SLARegQuery.GET_SLA_REG_ALL, jobId); 779 if (slaRegBean != null) { // filter out jobs picked by SLA job event listener 780 // but not actually configured for SLA 781 SLASummaryBean slaSummaryBean = SLASummaryQueryExecutor.getInstance().get( 782 SLASummaryQuery.GET_SLA_SUMMARY, jobId); 783 if (slaSummaryBean.getEventProcessed() < 7) { 784 slaCalc = new SLACalcStatus(slaSummaryBean, slaRegBean); 785 putAndIncrement(jobId, slaCalc); 786 } 787 } 788 } 789 } 790 if (slaCalc != null) { 791 synchronized (slaCalc) { 792 try { 793 // only get ZK lock when multiple servers running 794 boolean locked = true; 795 slaCalc.acquireLock(); 796 locked = slaCalc.isLocked(); 797 if (locked) { 798 // get eventProcessed on DB for validation in HA 799 Object eventProcObj = ((SLASummaryQueryExecutor) SLASummaryQueryExecutor.getInstance()) 800 .getSingleValue(SLASummaryQuery.GET_SLA_SUMMARY_EVENTPROCESSED, jobId); 801 byte eventProc = ((Byte) eventProcObj).byteValue(); 802 slaCalc.setEventProcessed(eventProc); 803 slaCalc.setJobStatus(jobStatus); 804 switch (jobEventStatus) { 805 case STARTED: 806 slaInfo = processJobStartSLA(slaCalc, startTime); 807 break; 808 case SUCCESS: 809 slaInfo = processJobEndSuccessSLA(slaCalc, startTime, endTime); 810 break; 811 case FAILURE: 812 slaInfo = processJobEndFailureSLA(slaCalc, startTime, endTime); 813 break; 814 default: 815 LOG.debug("Unknown Job Status for SLA purpose[{0}]", jobEventStatus); 816 slaInfo = getSLASummaryBean(slaCalc); 817 } 818 if (slaCalc.getEventProcessed() == 7) { 819 slaInfo.setEventProcessed(8); 820 slaMap.remove(jobId); 821 } 822 slaInfo.setLastModifiedTime(new Date()); 823 SLASummaryQueryExecutor.getInstance().executeUpdate( 824 SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_STATUS_ACTUAL_TIMES, slaInfo); 825 hasSla = true; 826 } 827 } 828 catch (InterruptedException e) { 829 throw new ServiceException(ErrorCode.E0606, slaCalc.getEntityKey(), slaCalc.getLockTimeOut()); 830 } 831 finally { 832 slaCalc.releaseLock(); 833 } 834 } 835 LOG.trace("SLA Status Event - Job:" + jobId + " Status:" + slaCalc.getSLAStatus()); 836 } 837 838 return hasSla; 839 } 840 841 /** 842 * Process SLA for jobs that started running. Also update actual-start time 843 * 844 * @param slaCalc 845 * @param actualStart 846 * @return SLASummaryBean 847 */ 848 private SLASummaryBean processJobStartSLA(SLACalcStatus slaCalc, Date actualStart) { 849 slaCalc.setActualStart(actualStart); 850 if (slaCalc.getSLAStatus().equals(SLAStatus.NOT_STARTED)) { 851 slaCalc.setSLAStatus(SLAStatus.IN_PROCESS); 852 } 853 SLARegistrationBean reg = slaCalc.getSLARegistrationBean(); 854 Date expecStart = reg.getExpectedStart(); 855 byte eventProc = slaCalc.getEventProcessed(); 856 // set event proc here 857 if (((eventProc & 1) == 0)) { 858 if (expecStart != null) { 859 if (actualStart.getTime() > expecStart.getTime()) { 860 slaCalc.setEventStatus(EventStatus.START_MISS); 861 } 862 else { 863 slaCalc.setEventStatus(EventStatus.START_MET); 864 } 865 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 866 } 867 eventProc += 1; 868 slaCalc.setEventProcessed(eventProc); 869 } 870 return getSLASummaryBean(slaCalc); 871 } 872 873 /** 874 * Process SLA for jobs that ended successfully. Also update actual-start 875 * and end time 876 * 877 * @param slaCalc 878 * @param actualStart 879 * @param actualEnd 880 * @return SLASummaryBean 881 * @throws JPAExecutorException 882 */ 883 private SLASummaryBean processJobEndSuccessSLA(SLACalcStatus slaCalc, Date actualStart, Date actualEnd) throws JPAExecutorException { 884 SLARegistrationBean reg = slaCalc.getSLARegistrationBean(); 885 slaCalc.setActualStart(actualStart); 886 slaCalc.setActualEnd(actualEnd); 887 long expectedDuration = reg.getExpectedDuration(); 888 long actualDuration = actualEnd.getTime() - actualStart.getTime(); 889 slaCalc.setActualDuration(actualDuration); 890 //check event proc 891 byte eventProc = slaCalc.getEventProcessed(); 892 if (((eventProc >> 1) & 1) == 0) { 893 processDurationSLA(expectedDuration, actualDuration, slaCalc); 894 eventProc += 2; 895 slaCalc.setEventProcessed(eventProc); 896 } 897 898 if (eventProc < 4) { 899 Date expectedEnd = reg.getExpectedEnd(); 900 if (actualEnd.getTime() > expectedEnd.getTime()) { 901 slaCalc.setEventStatus(EventStatus.END_MISS); 902 slaCalc.setSLAStatus(SLAStatus.MISS); 903 } 904 else { 905 slaCalc.setEventStatus(EventStatus.END_MET); 906 slaCalc.setSLAStatus(SLAStatus.MET); 907 } 908 eventProc += 4; 909 slaCalc.setEventProcessed(eventProc); 910 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 911 } 912 return getSLASummaryBean(slaCalc); 913 } 914 915 /** 916 * Process SLA for jobs that ended in failure. Also update actual-start and 917 * end time 918 * 919 * @param slaCalc 920 * @param actualStart 921 * @param actualEnd 922 * @return SLASummaryBean 923 * @throws JPAExecutorException 924 */ 925 private SLASummaryBean processJobEndFailureSLA(SLACalcStatus slaCalc, Date actualStart, Date actualEnd) throws JPAExecutorException { 926 slaCalc.setActualStart(actualStart); 927 slaCalc.setActualEnd(actualEnd); 928 if (actualStart == null) { // job failed before starting 929 if (slaCalc.getEventProcessed() < 4) { 930 slaCalc.setEventStatus(EventStatus.END_MISS); 931 slaCalc.setSLAStatus(SLAStatus.MISS); 932 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 933 slaCalc.setEventProcessed(7); 934 return getSLASummaryBean(slaCalc); 935 } 936 } 937 SLARegistrationBean reg = slaCalc.getSLARegistrationBean(); 938 long expectedDuration = reg.getExpectedDuration(); 939 long actualDuration = actualEnd.getTime() - actualStart.getTime(); 940 slaCalc.setActualDuration(actualDuration); 941 942 byte eventProc = slaCalc.getEventProcessed(); 943 if (((eventProc >> 1) & 1) == 0) { 944 if (expectedDuration != -1) { 945 slaCalc.setEventStatus(EventStatus.DURATION_MISS); 946 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 947 } 948 eventProc += 2; 949 slaCalc.setEventProcessed(eventProc); 950 } 951 if (eventProc < 4) { 952 slaCalc.setEventStatus(EventStatus.END_MISS); 953 slaCalc.setSLAStatus(SLAStatus.MISS); 954 eventProc += 4; 955 slaCalc.setEventProcessed(eventProc); 956 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 957 } 958 return getSLASummaryBean(slaCalc); 959 } 960 961 private SLASummaryBean getSLASummaryBean (SLACalcStatus slaCalc) { 962 SLASummaryBean slaSummaryBean = new SLASummaryBean(); 963 slaSummaryBean.setActualStart(slaCalc.getActualStart()); 964 slaSummaryBean.setActualEnd(slaCalc.getActualEnd()); 965 slaSummaryBean.setActualDuration(slaCalc.getActualDuration()); 966 slaSummaryBean.setSLAStatus(slaCalc.getSLAStatus()); 967 slaSummaryBean.setEventStatus(slaCalc.getEventStatus()); 968 slaSummaryBean.setEventProcessed(slaCalc.getEventProcessed()); 969 slaSummaryBean.setId(slaCalc.getId()); 970 slaSummaryBean.setJobStatus(slaCalc.getJobStatus()); 971 return slaSummaryBean; 972 } 973 974 private void processDurationSLA(long expected, long actual, SLACalcStatus slaCalc) { 975 if (expected != -1 && actual > expected) { 976 slaCalc.setEventStatus(EventStatus.DURATION_MISS); 977 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 978 } 979 else if (expected != -1 && actual <= expected) { 980 slaCalc.setEventStatus(EventStatus.DURATION_MET); 981 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 982 } 983 } 984 985 /* 986 * Confirm alerts against source of truth - DB. Also required in case of High Availability 987 */ 988 private void confirmWithDB(SLACalcStatus slaCalc) { 989 boolean ended = false, isEndMiss = false; 990 try { 991 switch (slaCalc.getAppType()) { 992 case WORKFLOW_JOB: 993 WorkflowJobBean wf = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(slaCalc.getId())); 994 if (wf.getEndTime() != null) { 995 ended = true; 996 if (wf.getStatus() == WorkflowJob.Status.KILLED || wf.getStatus() == WorkflowJob.Status.FAILED 997 || wf.getEndTime().getTime() > slaCalc.getExpectedEnd().getTime()) { 998 isEndMiss = true; 999 } 1000 } 1001 slaCalc.setActualStart(wf.getStartTime()); 1002 slaCalc.setActualEnd(wf.getEndTime()); 1003 slaCalc.setJobStatus(wf.getStatusStr()); 1004 break; 1005 case WORKFLOW_ACTION: 1006 WorkflowActionBean wa = jpaService.execute(new WorkflowActionGetForSLAJPAExecutor(slaCalc.getId())); 1007 if (wa.getEndTime() != null) { 1008 ended = true; 1009 if (wa.isTerminalWithFailure() 1010 || wa.getEndTime().getTime() > slaCalc.getExpectedEnd().getTime()) { 1011 isEndMiss = true; 1012 } 1013 } 1014 slaCalc.setActualStart(wa.getStartTime()); 1015 slaCalc.setActualEnd(wa.getEndTime()); 1016 slaCalc.setJobStatus(wa.getStatusStr()); 1017 break; 1018 case COORDINATOR_ACTION: 1019 CoordinatorActionBean ca = jpaService.execute(new CoordActionGetForSLAJPAExecutor(slaCalc.getId())); 1020 if (ca.isTerminalWithFailure()) { 1021 isEndMiss = ended = true; 1022 slaCalc.setActualEnd(ca.getLastModifiedTime()); 1023 } 1024 if (ca.getExternalId() != null) { 1025 wf = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(ca.getExternalId())); 1026 if (wf.getEndTime() != null) { 1027 ended = true; 1028 if (wf.getEndTime().getTime() > slaCalc.getExpectedEnd().getTime()) { 1029 isEndMiss = true; 1030 } 1031 } 1032 slaCalc.setActualEnd(wf.getEndTime()); 1033 slaCalc.setActualStart(wf.getStartTime()); 1034 } 1035 slaCalc.setJobStatus(ca.getStatusStr()); 1036 break; 1037 default: 1038 LOG.debug("Unsupported App-type for SLA - " + slaCalc.getAppType()); 1039 } 1040 1041 byte eventProc = slaCalc.getEventProcessed(); 1042 if (ended) { 1043 if (isEndMiss) { 1044 slaCalc.setSLAStatus(SLAStatus.MISS); 1045 } 1046 else { 1047 slaCalc.setSLAStatus(SLAStatus.MET); 1048 } 1049 if (slaCalc.getActualStart() != null) { 1050 if ((eventProc & 1) == 0) { 1051 if (slaCalc.getExpectedStart().getTime() < slaCalc.getActualStart().getTime()) { 1052 slaCalc.setEventStatus(EventStatus.START_MISS); 1053 } 1054 else { 1055 slaCalc.setEventStatus(EventStatus.START_MET); 1056 } 1057 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 1058 } 1059 slaCalc.setActualDuration(slaCalc.getActualEnd().getTime() - slaCalc.getActualStart().getTime()); 1060 if (((eventProc >> 1) & 1) == 0) { 1061 processDurationSLA(slaCalc.getExpectedDuration(), slaCalc.getActualDuration(), slaCalc); 1062 } 1063 } 1064 if (eventProc < 4) { 1065 if (isEndMiss) { 1066 slaCalc.setEventStatus(EventStatus.END_MISS); 1067 } 1068 else { 1069 slaCalc.setEventStatus(EventStatus.END_MET); 1070 } 1071 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 1072 } 1073 slaCalc.setEventProcessed(8); 1074 } 1075 else { 1076 if (slaCalc.getActualStart() != null) { 1077 slaCalc.setSLAStatus(SLAStatus.IN_PROCESS); 1078 } 1079 if ((eventProc & 1) == 0) { 1080 if (slaCalc.getActualStart() != null) { 1081 if (slaCalc.getExpectedStart().getTime() < slaCalc.getActualStart().getTime()) { 1082 slaCalc.setEventStatus(EventStatus.START_MISS); 1083 } 1084 else { 1085 slaCalc.setEventStatus(EventStatus.START_MET); 1086 } 1087 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 1088 eventProc++; 1089 } 1090 else if (slaCalc.getExpectedStart().getTime() < System.currentTimeMillis()) { 1091 slaCalc.setEventStatus(EventStatus.START_MISS); 1092 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 1093 eventProc++; 1094 } 1095 } 1096 if (((eventProc >> 1) & 1) == 0 && slaCalc.getActualStart() != null 1097 && slaCalc.getExpectedDuration() != -1) { 1098 if (System.currentTimeMillis() - slaCalc.getActualStart().getTime() > slaCalc.getExpectedDuration()) { 1099 slaCalc.setEventStatus(EventStatus.DURATION_MISS); 1100 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 1101 eventProc += 2; 1102 } 1103 } 1104 if (eventProc < 4 && slaCalc.getExpectedEnd().getTime() < System.currentTimeMillis()) { 1105 slaCalc.setEventStatus(EventStatus.END_MISS); 1106 slaCalc.setSLAStatus(SLAStatus.MISS); 1107 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 1108 eventProc += 4; 1109 } 1110 slaCalc.setEventProcessed(eventProc); 1111 } 1112 } 1113 catch (Exception e) { 1114 LOG.warn("Error while confirming SLA against DB for jobid= " + slaCalc.getId() + ". Exception is " 1115 + e.getClass().getName() + ": " + e.getMessage()); 1116 if (slaCalc.getEventProcessed() < 4 && slaCalc.getExpectedEnd().getTime() < System.currentTimeMillis()) { 1117 slaCalc.setEventStatus(EventStatus.END_MISS); 1118 slaCalc.setSLAStatus(SLAStatus.MISS); 1119 eventHandler.queueEvent(new SLACalcStatus(slaCalc)); 1120 slaCalc.setEventProcessed(slaCalc.getEventProcessed() + 4); 1121 } 1122 } 1123 } 1124 1125 @VisibleForTesting 1126 public boolean isJobIdInSLAMap(String jobId) { 1127 return this.slaMap.containsKey(jobId); 1128 } 1129 1130 @VisibleForTesting 1131 public boolean isJobIdInHistorySet(String jobId) { 1132 return this.historySet.contains(jobId); 1133 } 1134 1135 private void setLogPrefix(String jobId) { 1136 LOG = LogUtils.setLogInfo(LOG, jobId, null, null); 1137 } 1138 1139 private boolean putAndIncrement(final String jobId, final SLACalcStatus newStatus) { 1140 if (slaMap.put(jobId, newStatus) == null) { 1141 LOG.trace("Added a new item to SLA map. [jobId={0}]", jobId); 1142 instrumentation.incr(INSTRUMENTATION_GROUP, SLA_MAP, 1); 1143 return true; 1144 } 1145 1146 LOG.trace("Updated an existing item in SLA map. [jobId={0}]", jobId); 1147 return false; 1148 } 1149 1150 private boolean removeAndDecrement(final String jobId) { 1151 if (slaMap.remove(jobId) != null) { 1152 LOG.trace("Removed an existing item from SLA map. [jobId={0}]", jobId); 1153 instrumentation.decr(INSTRUMENTATION_GROUP, SLA_MAP, 1); 1154 return true; 1155 } 1156 1157 LOG.trace("Tried to remove a non-existing item from SLA map. [jobId={0}]", jobId); 1158 return false; 1159 } 1160 1161 private void resetRetryCount(final String jobId) { 1162 if (slaMap.containsKey(jobId)) { 1163 LOG.debug("Resetting retry count on [{0}]", jobId); 1164 final SLACalcStatus existingStatus = slaMap.get(jobId); 1165 existingStatus.resetRetryCount(); 1166 putAndIncrement(jobId, existingStatus); 1167 } 1168 } 1169 1170 private void incrementRetryCountAndRemove(final String jobId) { 1171 LOG.debug("Checking SLA calculator status [{0}] for retry count", jobId); 1172 if (slaMap.containsKey(jobId)) { 1173 final SLACalcStatus existingStatus = slaMap.get(jobId); 1174 if (existingStatus.getRetryCount() < maxRetryCount) { 1175 existingStatus.incrementRetryCount(); 1176 LOG.debug("Retrying with SLA calculator status [{0}] retry count [{1}]", jobId, existingStatus.getRetryCount()); 1177 putAndIncrement(jobId, existingStatus); 1178 } 1179 else { 1180 LOG.debug("Removing [{0}] from SLA map as maximum retry count reached", jobId); 1181 removeAndDecrement(jobId); 1182 } 1183 } 1184 } 1185}