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.action.hadoop;
020
021import org.apache.hadoop.conf.Configuration;
022import org.apache.hadoop.fs.Path;
023import org.apache.hadoop.mapred.Counters;
024import org.apache.hadoop.mapred.JobClient;
025import org.apache.hadoop.mapred.JobConf;
026import org.apache.hadoop.mapred.JobID;
027import org.apache.hadoop.mapred.RunningJob;
028import org.apache.oozie.action.ActionExecutorException;
029import org.apache.oozie.client.WorkflowAction;
030import org.apache.oozie.util.XConfiguration;
031import org.apache.oozie.util.XLog;
032import org.apache.oozie.util.XmlUtils;
033import org.jdom.Element;
034import org.jdom.Namespace;
035
036import java.io.IOException;
037import java.io.StringReader;
038import java.util.ArrayList;
039import java.util.List;
040import java.util.StringTokenizer;
041
042public class SqoopActionExecutor extends JavaActionExecutor {
043
044  public static final String OOZIE_ACTION_EXTERNAL_STATS_WRITE = "oozie.action.external.stats.write";
045  private static final String SQOOP_MAIN_CLASS_NAME = "org.apache.oozie.action.hadoop.SqoopMain";
046  static final String SQOOP_ARGS = "oozie.sqoop.args";
047  private static final String SQOOP = "sqoop";
048
049    public SqoopActionExecutor() {
050        super(SQOOP);
051    }
052
053    @Override
054    public List<Class> getLauncherClasses() {
055        List<Class> classes = new ArrayList<Class>();
056        try {
057            classes.add(Class.forName(SQOOP_MAIN_CLASS_NAME));
058        }
059        catch (ClassNotFoundException e) {
060            throw new RuntimeException("Class not found", e);
061        }
062        return classes;
063    }
064
065    @Override
066    protected String getLauncherMain(Configuration launcherConf, Element actionXml) {
067        return launcherConf.get(LauncherMapper.CONF_OOZIE_ACTION_MAIN_CLASS, SQOOP_MAIN_CLASS_NAME);
068    }
069
070    @Override
071    @SuppressWarnings("unchecked")
072    Configuration setupActionConf(Configuration actionConf, Context context, Element actionXml, Path appPath)
073            throws ActionExecutorException {
074        super.setupActionConf(actionConf, context, actionXml, appPath);
075        Namespace ns = actionXml.getNamespace();
076
077        try {
078            Element e = actionXml.getChild("configuration", ns);
079            if (e != null) {
080                String strConf = XmlUtils.prettyPrint(e).toString();
081                XConfiguration inlineConf = new XConfiguration(new StringReader(strConf));
082                checkForDisallowedProps(inlineConf, "inline configuration");
083                XConfiguration.copy(inlineConf, actionConf);
084            }
085        } catch (IOException ex) {
086            throw convertException(ex);
087        }
088
089        final List<String> argList = new ArrayList<>();
090        // Build a list of arguments from either a tokenized <command> string or a list of <arg>
091        if (actionXml.getChild("command", ns) != null) {
092            String command = actionXml.getChild("command", ns).getTextTrim();
093            StringTokenizer st = new StringTokenizer(command, " ");
094            while (st.hasMoreTokens()) {
095                argList.add(st.nextToken());
096            }
097        }
098        else {
099            List<Element> eArgs = (List<Element>) actionXml.getChildren("arg", ns);
100            for (Element elem : eArgs) {
101                argList.add(elem.getTextTrim());
102            }
103        }
104        // If the command is given accidentally as "sqoop import --option"
105        // instead of "import --option" we can make a user's life easier
106        // by removing away the unnecessary "sqoop" token.
107        // However, we do not do this if the command looks like
108        // "sqoop --option", as that's entirely invalid.
109        if (argList.size() > 1 &&
110                argList.get(0).equalsIgnoreCase(SQOOP) &&
111                !argList.get(1).startsWith("-")) {
112            XLog.getLog(getClass()).info(
113                    "Found a redundant 'sqoop' prefixing the command. Removing it.");
114            argList.remove(0);
115        }
116
117        setSqoopCommand(actionConf, argList.toArray(new String[argList.size()]));
118        return actionConf;
119    }
120
121    private void setSqoopCommand(Configuration conf, String[] args) {
122        MapReduceMain.setStrings(conf, SQOOP_ARGS, args);
123    }
124
125    /**
126     * We will gather counters from all executed action Hadoop jobs (e.g. jobs
127     * that moved data, not the launcher itself) and merge them together. There
128     * will be only one job most of the time. The only exception is
129     * import-all-table option that will execute one job per one exported table.
130     *
131     * @param context Action context
132     * @param action Workflow action
133     * @throws ActionExecutorException
134     */
135    @Override
136    public void end(Context context, WorkflowAction action) throws ActionExecutorException {
137        super.end(context, action);
138        JobClient jobClient = null;
139
140        boolean exception = false;
141        try {
142            if (action.getStatus() == WorkflowAction.Status.OK) {
143                Element actionXml = XmlUtils.parseXml(action.getConf());
144                JobConf jobConf = createBaseHadoopConf(context, actionXml);
145                jobClient = createJobClient(context, jobConf);
146
147                // Cumulative counters for all Sqoop mapreduce jobs
148                Counters counters = null;
149
150                // Sqoop do not have to create mapreduce job each time
151                String externalIds = action.getExternalChildIDs();
152                if (externalIds != null && !externalIds.trim().isEmpty()) {
153                    String []jobIds = externalIds.split(",");
154
155                    for(String jobId : jobIds) {
156                        RunningJob runningJob = jobClient.getJob(JobID.forName(jobId));
157                        if (runningJob == null) {
158                          throw new ActionExecutorException(ActionExecutorException.ErrorType.FAILED, "SQOOP001",
159                            "Unknown hadoop job [{0}] associated with action [{1}].  Failing this action!", action
160                            .getExternalId(), action.getId());
161                        }
162
163                        Counters taskCounters = runningJob.getCounters();
164                        if(taskCounters != null) {
165                            if(counters == null) {
166                              counters = taskCounters;
167                            } else {
168                              counters.incrAllCounters(taskCounters);
169                            }
170                        } else {
171                          XLog.getLog(getClass()).warn("Could not find Hadoop Counters for job: [{0}]", jobId);
172                        }
173                    }
174                }
175
176                if (counters != null) {
177                    ActionStats stats = new MRStats(counters);
178                    String statsJsonString = stats.toJSON();
179                    context.setVar(MapReduceActionExecutor.HADOOP_COUNTERS, statsJsonString);
180
181                    // If action stats write property is set to false by user or
182                    // size of stats is greater than the maximum allowed size,
183                    // do not store the action stats
184                    if (Boolean.parseBoolean(evaluateConfigurationProperty(actionXml,
185                            OOZIE_ACTION_EXTERNAL_STATS_WRITE, "true"))
186                            && (statsJsonString.getBytes().length <= getMaxExternalStatsSize())) {
187                        context.setExecutionStats(statsJsonString);
188                        LOG.debug(
189                          "Printing stats for sqoop action as a JSON string : [{0}]", statsJsonString);
190                    }
191                } else {
192                    context.setVar(MapReduceActionExecutor.HADOOP_COUNTERS, "");
193                    XLog.getLog(getClass()).warn("Can't find any associated Hadoop job counters");
194                }
195            }
196        }
197        catch (Exception ex) {
198            exception = true;
199            throw convertException(ex);
200        }
201        finally {
202            if (jobClient != null) {
203                try {
204                    jobClient.close();
205                }
206                catch (Exception e) {
207                    if (exception) {
208                        LOG.error("JobClient error: ", e);
209                    }
210                    else {
211                        throw convertException(e);
212                    }
213                }
214            }
215        }
216    }
217
218    // Return the value of the specified configuration property
219    private String evaluateConfigurationProperty(Element actionConf, String key, String defaultValue)
220            throws ActionExecutorException {
221        try {
222            if (actionConf != null) {
223                Namespace ns = actionConf.getNamespace();
224                Element e = actionConf.getChild("configuration", ns);
225
226                if(e != null) {
227                  String strConf = XmlUtils.prettyPrint(e).toString();
228                  XConfiguration inlineConf = new XConfiguration(new StringReader(strConf));
229                  return inlineConf.get(key, defaultValue);
230                }
231            }
232            return defaultValue;
233        }
234        catch (IOException ex) {
235            throw convertException(ex);
236        }
237    }
238
239    /**
240     * Return the sharelib name for the action.
241     *
242     * @return returns <code>sqoop</code>.
243     * @param actionXml
244     */
245    @Override
246    protected String getDefaultShareLibName(Element actionXml) {
247        return "sqoop";
248    }
249
250}