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.bundle;
020
021import java.io.IOException;
022import java.io.StringReader;
023import java.util.Date;
024import java.util.HashMap;
025import java.util.List;
026import java.util.Map;
027import java.util.Map.Entry;
028
029import org.apache.hadoop.conf.Configuration;
030import org.apache.oozie.BundleActionBean;
031import org.apache.oozie.BundleJobBean;
032import org.apache.oozie.ErrorCode;
033import org.apache.oozie.XException;
034import org.apache.oozie.action.hadoop.OozieJobInfo;
035import org.apache.oozie.client.Job;
036import org.apache.oozie.client.OozieClient;
037import org.apache.oozie.client.rest.JsonBean;
038import org.apache.oozie.command.CommandException;
039import org.apache.oozie.command.PreconditionException;
040import org.apache.oozie.command.StartTransitionXCommand;
041import org.apache.oozie.command.coord.CoordSubmitXCommand;
042import org.apache.oozie.executor.jpa.BatchQueryExecutor;
043import org.apache.oozie.executor.jpa.BundleJobQueryExecutor;
044import org.apache.oozie.executor.jpa.BundleJobQueryExecutor.BundleJobQuery;
045import org.apache.oozie.executor.jpa.JPAExecutorException;
046import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry;
047import org.apache.oozie.util.ConfigUtils;
048import org.apache.oozie.util.ELUtils;
049import org.apache.oozie.util.JobUtils;
050import org.apache.oozie.util.LogUtils;
051import org.apache.oozie.util.ParamChecker;
052import org.apache.oozie.util.XConfiguration;
053import org.apache.oozie.util.XmlUtils;
054import org.jdom.Attribute;
055import org.jdom.Element;
056import org.jdom.JDOMException;
057
058/**
059 * The command to start Bundle job
060 */
061public class BundleStartXCommand extends StartTransitionXCommand {
062    private final String jobId;
063    private BundleJobBean bundleJob;
064
065    /**
066     * The constructor for class {@link BundleStartXCommand}
067     *
068     * @param jobId the bundle job id
069     */
070    public BundleStartXCommand(String jobId) {
071        super("bundle_start", "bundle_start", 1);
072        this.jobId = ParamChecker.notEmpty(jobId, "jobId");
073    }
074
075    /**
076     * The constructor for class {@link BundleStartXCommand}
077     *
078     * @param jobId the bundle job id
079     * @param dryrun true if dryrun is enable
080     */
081    public BundleStartXCommand(String jobId, boolean dryrun) {
082        super("bundle_start", "bundle_start", 1, dryrun);
083        this.jobId = ParamChecker.notEmpty(jobId, "jobId");
084    }
085
086    /* (non-Javadoc)
087     * @see org.apache.oozie.command.XCommand#getEntityKey()
088     */
089    @Override
090    public String getEntityKey() {
091        return jobId;
092    }
093
094    @Override
095    public String getKey() {
096        return getName() + "_" + jobId;
097    }
098
099    /* (non-Javadoc)
100     * @see org.apache.oozie.command.XCommand#isLockRequired()
101     */
102    @Override
103    protected boolean isLockRequired() {
104        return true;
105    }
106
107    @Override
108    protected void verifyPrecondition() throws CommandException, PreconditionException {
109        if (bundleJob.getStatus() != Job.Status.PREP) {
110            String msg = "Bundle " + bundleJob.getId() + " is not in PREP status. It is in : " + bundleJob.getStatus();
111            LOG.info(msg);
112            throw new PreconditionException(ErrorCode.E1100, msg);
113        }
114    }
115
116    @Override
117    public void loadState() throws CommandException {
118        try {
119            this.bundleJob = BundleJobQueryExecutor.getInstance().get(BundleJobQuery.GET_BUNDLE_JOB, jobId);
120            LogUtils.setLogInfo(bundleJob);
121            super.setJob(bundleJob);
122        }
123        catch (XException ex) {
124            throw new CommandException(ex);
125        }
126    }
127
128    /* (non-Javadoc)
129     * @see org.apache.oozie.command.StartTransitionXCommand#StartChildren()
130     */
131    @Override
132    public void StartChildren() throws CommandException {
133        LOG.debug("Started coord jobs for the bundle=[{0}]", jobId);
134        insertBundleActions();
135        startCoordJobs();
136        LOG.debug("Ended coord jobs for the bundle=[{0}]", jobId);
137    }
138
139    /* (non-Javadoc)
140     * @see org.apache.oozie.command.TransitionXCommand#notifyParent()
141     */
142    @Override
143    public void notifyParent() {
144    }
145
146    /* (non-Javadoc)
147     * @see org.apache.oozie.command.StartTransitionXCommand#performWrites()
148     */
149    @Override
150    public void performWrites() throws CommandException {
151        try {
152            BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, updateList, null);
153        }
154        catch (JPAExecutorException e) {
155            throw new CommandException(e);
156        }
157    }
158
159    /**
160     * Insert bundle actions
161     *
162     * @throws CommandException thrown if failed to create bundle actions
163     */
164    @SuppressWarnings("unchecked")
165    private void insertBundleActions() throws CommandException {
166        if (bundleJob != null) {
167            Map<String, Boolean> map = new HashMap<String, Boolean>();
168            try {
169                Element bAppXml = XmlUtils.parseXml(bundleJob.getJobXml());
170                List<Element> coordElems = bAppXml.getChildren("coordinator", bAppXml.getNamespace());
171                for (Element elem : coordElems) {
172                    Attribute name = elem.getAttribute("name");
173                    Attribute critical = elem.getAttribute("critical");
174                    if (name != null) {
175                        if (map.containsKey(name.getValue())) {
176                            throw new CommandException(ErrorCode.E1304, name);
177                        }
178                        boolean isCritical = false;
179                        if (critical != null && Boolean.parseBoolean(critical.getValue())) {
180                            isCritical = true;
181                        }
182                        map.put(name.getValue(), isCritical);
183                    }
184                    else {
185                        throw new CommandException(ErrorCode.E1305);
186                    }
187                }
188            }
189            catch (JDOMException jex) {
190                throw new CommandException(ErrorCode.E1301, jex.getMessage(), jex);
191            }
192
193            // if there is no coordinator for this bundle, failed it.
194            if (map.isEmpty()) {
195                bundleJob.setStatus(Job.Status.FAILED);
196                bundleJob.resetPending();
197                try {
198                    BundleJobQueryExecutor.getInstance().executeUpdate(BundleJobQuery.UPDATE_BUNDLE_JOB_STATUS_PENDING, bundleJob);
199                }
200                catch (JPAExecutorException jex) {
201                    throw new CommandException(jex);
202                }
203
204                LOG.debug("No coord jobs for the bundle=[{0}], failed it!!", jobId);
205                throw new CommandException(ErrorCode.E1318, jobId);
206            }
207
208            for (Entry<String, Boolean> coordName : map.entrySet()) {
209                BundleActionBean action = createBundleAction(jobId, coordName.getKey(), coordName.getValue());
210                insertList.add(action);
211            }
212        }
213        else {
214            throw new CommandException(ErrorCode.E0604, jobId);
215        }
216    }
217
218    private BundleActionBean createBundleAction(String jobId, String coordName, boolean isCritical) {
219        BundleActionBean action = new BundleActionBean();
220        action.setBundleActionId(jobId + "_" + coordName);
221        action.setBundleId(jobId);
222        action.setCoordName(coordName);
223        action.setStatus(Job.Status.PREP);
224        action.setLastModifiedTime(new Date());
225        if (isCritical) {
226            action.setCritical();
227        }
228        else {
229            action.resetCritical();
230        }
231        return action;
232    }
233
234    /**
235     * Start Coord Jobs
236     *
237     * @throws CommandException thrown if failed to start coord jobs
238     */
239    @SuppressWarnings("unchecked")
240    private void startCoordJobs() throws CommandException {
241        if (bundleJob != null) {
242            try {
243                Element bAppXml = XmlUtils.parseXml(bundleJob.getJobXml());
244                List<Element> coordElems = bAppXml.getChildren("coordinator", bAppXml.getNamespace());
245                for (Element coordElem : coordElems) {
246                    Attribute name = coordElem.getAttribute("name");
247
248                    Configuration coordConf = mergeConfig(coordElem);
249                    coordConf.set(OozieClient.BUNDLE_ID, jobId);
250                    if (OozieJobInfo.isJobInfoEnabled()) {
251                        coordConf.set(OozieJobInfo.BUNDLE_NAME, bundleJob.getAppName());
252                    }
253                    String coordName=name.getValue();
254                    try {
255                        coordName = ELUtils.resolveAppName(coordName, coordConf);
256                    }
257                    catch (Exception e) {
258                        throw new CommandException(ErrorCode.E1321, e.getMessage(), e);
259
260                    }
261                    queue(new CoordSubmitXCommand(coordConf, bundleJob.getId(), name.getValue()));
262
263                }
264                updateBundleAction();
265            }
266            catch (JDOMException jex) {
267                throw new CommandException(ErrorCode.E1301, jex.getMessage(), jex);
268            }
269            catch (JPAExecutorException je) {
270                throw new CommandException(je);
271            }
272        }
273        else {
274            throw new CommandException(ErrorCode.E0604, jobId);
275        }
276    }
277
278    private void updateBundleAction() throws JPAExecutorException {
279        for(JsonBean bAction : insertList) {
280            BundleActionBean action = (BundleActionBean) bAction;
281            action.incrementAndGetPending();
282            action.setLastModifiedTime(new Date());
283        }
284    }
285
286    /**
287     * Merge Bundle job config and the configuration from the coord job to pass
288     * to Coord Engine
289     *
290     * @param coordElem the coordinator configuration
291     * @return Configuration merged configuration
292     * @throws CommandException thrown if failed to merge configuration
293     */
294    private Configuration mergeConfig(Element coordElem) throws CommandException {
295        String jobConf = bundleJob.getConf();
296        // Step 1: runConf = jobConf
297        Configuration runConf = null;
298        try {
299            runConf = new XConfiguration(new StringReader(jobConf));
300        }
301        catch (IOException e1) {
302            LOG.warn("Configuration parse error in:" + jobConf);
303            throw new CommandException(ErrorCode.E1306, e1.getMessage(), e1);
304        }
305        // Step 2: Merge local properties into runConf
306        // extract 'property' tags under 'configuration' block in the coordElem
307        // convert Element to XConfiguration
308        Element localConfigElement = coordElem.getChild("configuration", coordElem.getNamespace());
309
310        if (localConfigElement != null) {
311            String strConfig = XmlUtils.prettyPrint(localConfigElement).toString();
312            Configuration localConf;
313            try {
314                localConf = new XConfiguration(new StringReader(strConfig));
315            }
316            catch (IOException e1) {
317                LOG.warn("Configuration parse error in:" + strConfig);
318                throw new CommandException(ErrorCode.E1307, e1.getMessage(), e1);
319            }
320
321            // copy configuration properties in the coordElem to the runConf
322            XConfiguration.copy(localConf, runConf);
323
324            ConfigUtils.checkAndSetDisallowedProperties(runConf,
325                    bundleJob.getUser(),
326                    new CommandException(ErrorCode.E1303,
327                            String.format("%s=%s", OozieClient.USER_NAME, runConf.get(OozieClient.USER_NAME)),
328                            bundleJob.getUser()),
329                    true);
330        }
331
332        // Step 3: Extract value of 'app-path' in coordElem, save it as a
333        // new property called 'oozie.coord.application.path', and normalize.
334        String appPath = coordElem.getChild("app-path", coordElem.getNamespace()).getValue();
335        runConf.set(OozieClient.COORDINATOR_APP_PATH, appPath);
336        // Normalize coordinator appPath here;
337        try {
338            JobUtils.normalizeAppPath(runConf.get(OozieClient.USER_NAME), runConf.get(OozieClient.GROUP_NAME), runConf);
339        }
340        catch (IOException e) {
341            throw new CommandException(ErrorCode.E1001, runConf.get(OozieClient.COORDINATOR_APP_PATH));
342        }
343        return runConf;
344    }
345
346    /* (non-Javadoc)
347     * @see org.apache.oozie.command.TransitionXCommand#getJob()
348     */
349    @Override
350    public Job getJob() {
351        return bundleJob;
352    }
353
354    /* (non-Javadoc)
355     * @see org.apache.oozie.command.TransitionXCommand#updateJob()
356     */
357    @Override
358    public void updateJob() throws CommandException {
359        updateList.add(new UpdateEntry<BundleJobQuery>(BundleJobQuery.UPDATE_BUNDLE_JOB_STATUS_PENDING, bundleJob));
360    }
361}