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.service; 020 021import java.util.Date; 022 023import org.apache.hadoop.conf.Configuration; 024import org.apache.oozie.ErrorCode; 025import org.apache.oozie.client.event.JobEvent.EventStatus; 026import org.apache.oozie.executor.jpa.JPAExecutorException; 027import org.apache.oozie.service.ConfigurationService; 028import org.apache.oozie.service.EventHandlerService; 029import org.apache.oozie.service.SchedulerService; 030import org.apache.oozie.service.Service; 031import org.apache.oozie.service.ServiceException; 032import org.apache.oozie.service.Services; 033import org.apache.oozie.sla.SLACalculator; 034import org.apache.oozie.sla.SLACalculatorMemory; 035import org.apache.oozie.sla.SLARegistrationBean; 036import org.apache.oozie.util.XLog; 037 038import com.google.common.annotations.VisibleForTesting; 039 040public class SLAService implements Service { 041 042 public static final String CONF_PREFIX = "oozie.sla.service.SLAService."; 043 public static final String CONF_CALCULATOR_IMPL = CONF_PREFIX + "calculator.impl"; 044 public static final String CONF_CAPACITY = CONF_PREFIX + "capacity"; 045 public static final String CONF_ALERT_EVENTS = CONF_PREFIX + "alert.events"; 046 public static final String CONF_EVENTS_MODIFIED_AFTER = CONF_PREFIX + "events.modified.after"; 047 public static final String CONF_JOB_EVENT_LATENCY = CONF_PREFIX + "job.event.latency"; 048 //Time interval, in seconds, at which SLA Worker will be scheduled to run 049 public static final String CONF_SLA_CHECK_INTERVAL = CONF_PREFIX + "check.interval"; 050 public static final String CONF_SLA_CHECK_INITIAL_DELAY = CONF_PREFIX + "check.initial.delay"; 051 public static final String CONF_SLA_CALC_LOCK_TIMEOUT = CONF_PREFIX + "oozie.sla.calc.default.lock.timeout"; 052 public static final String CONF_SLA_HISTORY_PURGE_INTERVAL = CONF_PREFIX + "history.purge.interval"; 053 public static final String CONF_MAXIMUM_RETRY_COUNT = CONF_PREFIX + "maximum.retry.count"; 054 055 private static SLACalculator calcImpl; 056 private static boolean slaEnabled = false; 057 private EventHandlerService eventHandler; 058 public static XLog LOG; 059 @Override 060 public void init(Services services) throws ServiceException { 061 try { 062 Configuration conf = services.getConf(); 063 Class<? extends SLACalculator> calcClazz = (Class<? extends SLACalculator>) ConfigurationService.getClass( 064 conf, CONF_CALCULATOR_IMPL); 065 calcImpl = calcClazz == null ? new SLACalculatorMemory() : (SLACalculator) calcClazz.newInstance(); 066 calcImpl.init(conf); 067 eventHandler = Services.get().get(EventHandlerService.class); 068 if (eventHandler == null) { 069 throw new ServiceException(ErrorCode.E0103, "EventHandlerService", "Add it under config " 070 + Services.CONF_SERVICE_EXT_CLASSES + " or declare it BEFORE SLAService"); 071 } 072 LOG = XLog.getLog(getClass()); 073 java.util.Set<String> appTypes = eventHandler.getAppTypes(); 074 appTypes.add("workflow_action"); 075 eventHandler.setAppTypes(appTypes); 076 077 Runnable slaThread = new SLAWorker(calcImpl); 078 // schedule runnable by default every 30 sec 079 int slaCheckInterval = ConfigurationService.getInt(conf, CONF_SLA_CHECK_INTERVAL); 080 int slaCheckInitialDelay = ConfigurationService.getInt(conf, CONF_SLA_CHECK_INITIAL_DELAY); 081 services.get(SchedulerService.class).schedule(slaThread, slaCheckInitialDelay, slaCheckInterval, 082 SchedulerService.Unit.SEC); 083 slaEnabled = true; 084 LOG.info("SLAService initialized with impl [{0}] capacity [{1}]", calcImpl.getClass().getName(), 085 conf.get(SLAService.CONF_CAPACITY)); 086 } 087 catch (Exception ex) { 088 throw new ServiceException(ErrorCode.E0102, ex.getMessage(), ex); 089 } 090 } 091 092 @Override 093 public void destroy() { 094 slaEnabled = false; 095 } 096 097 @Override 098 public Class<? extends Service> getInterface() { 099 return SLAService.class; 100 } 101 102 public static boolean isEnabled() { 103 return slaEnabled; 104 } 105 106 @VisibleForTesting 107 public SLACalculator getSLACalculator() { 108 return calcImpl; 109 } 110 111 @VisibleForTesting 112 public void runSLAWorker() { 113 new SLAWorker(calcImpl).run(); 114 } 115 116 private class SLAWorker implements Runnable { 117 118 SLACalculator calc; 119 120 public SLAWorker(SLACalculator calc) { 121 this.calc = calc; 122 } 123 124 @Override 125 public void run() { 126 if (Thread.currentThread().isInterrupted()) { 127 return; 128 } 129 try { 130 calc.updateAllSlaStatus(); 131 } 132 catch (Throwable error) { 133 XLog.getLog(SLAService.class).debug("Throwable in SLAWorker thread run : ", error); 134 } 135 } 136 } 137 138 public boolean addRegistrationEvent(SLARegistrationBean reg) throws ServiceException { 139 try { 140 if (calcImpl.addRegistration(reg.getId(), reg)) { 141 return true; 142 } 143 else { 144 LOG.warn("SLA queue full. Unable to add new SLA entry for job [{0}]", reg.getId()); 145 } 146 } 147 catch (JPAExecutorException ex) { 148 LOG.warn("Could not add new SLA entry for job [{0}]", reg.getId(), ex); 149 } 150 return false; 151 } 152 153 public boolean updateRegistrationEvent(SLARegistrationBean reg) throws ServiceException { 154 try { 155 if (calcImpl.updateRegistration(reg.getId(), reg)) { 156 return true; 157 } 158 else { 159 LOG.warn("SLA queue full. Unable to update the SLA entry for job [{0}]", reg.getId()); 160 } 161 } 162 catch (JPAExecutorException ex) { 163 LOG.warn("Could not update SLA entry for job [{0}]", reg.getId(), ex); 164 } 165 return false; 166 } 167 168 public boolean addStatusEvent(String jobId, String status, EventStatus eventStatus, Date startTime, Date endTime) 169 throws ServiceException { 170 try { 171 if (calcImpl.addJobStatus(jobId, status, eventStatus, startTime, endTime)) { 172 return true; 173 } 174 } 175 catch (JPAExecutorException jpe) { 176 LOG.error("Exception while adding SLA Status event for Job [{0}]", jobId); 177 } 178 return false; 179 } 180 181 public void removeRegistration(String jobId) { 182 calcImpl.removeRegistration(jobId); 183 } 184 185}