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}