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 java.io.BufferedReader;
022import java.io.IOException;
023import java.io.InputStream;
024import java.io.InputStreamReader;
025import java.io.OutputStream;
026import java.math.BigInteger;
027import java.security.MessageDigest;
028import java.security.NoSuchAlgorithmException;
029import java.security.PrivilegedExceptionAction;
030import java.util.ArrayList;
031import java.util.Collection;
032import java.util.HashMap;
033import java.util.List;
034import java.util.Map;
035import java.util.Properties;
036
037import org.apache.hadoop.conf.Configuration;
038import org.apache.hadoop.fs.FileSystem;
039import org.apache.hadoop.fs.Path;
040import org.apache.hadoop.io.SequenceFile;
041import org.apache.hadoop.io.Text;
042import org.apache.hadoop.mapred.JobConf;
043import org.apache.hadoop.mapred.RunningJob;
044import org.apache.hadoop.mapred.Counters;
045import org.apache.hadoop.security.UserGroupInformation;
046import org.apache.oozie.client.OozieClient;
047import org.apache.oozie.client.WorkflowAction;
048import org.apache.oozie.service.HadoopAccessorException;
049import org.apache.oozie.service.HadoopAccessorService;
050import org.apache.oozie.service.Services;
051import org.apache.oozie.service.URIHandlerService;
052import org.apache.oozie.service.UserGroupInformationService;
053import org.apache.oozie.util.IOUtils;
054import org.apache.oozie.util.PropertiesUtils;
055
056public class LauncherMapperHelper {
057
058    public static final String OOZIE_ACTION_YARN_TAG = "oozie.action.yarn.tag";
059
060    private static final LauncherInputFormatClassLocator launcherInputFormatClassLocator = new LauncherInputFormatClassLocator();
061
062    public static String getRecoveryId(Configuration launcherConf, Path actionDir, String recoveryId)
063            throws HadoopAccessorException, IOException {
064        String jobId = null;
065        Path recoveryFile = new Path(actionDir, recoveryId);
066        FileSystem fs = Services.get().get(HadoopAccessorService.class)
067                .createFileSystem(launcherConf.get("user.name"),recoveryFile.toUri(), launcherConf);
068
069        if (fs.exists(recoveryFile)) {
070            InputStream is = fs.open(recoveryFile);
071            BufferedReader reader = new BufferedReader(new InputStreamReader(is));
072            jobId = reader.readLine();
073            reader.close();
074        }
075        return jobId;
076
077    }
078
079    public static void setupMainClass(Configuration launcherConf, String javaMainClass) {
080        // Only set the javaMainClass if its not null or empty string, this way the user can override the action's main class via
081        // <configuration> property
082        if (javaMainClass != null && !javaMainClass.equals("")) {
083            launcherConf.set(LauncherMapper.CONF_OOZIE_ACTION_MAIN_CLASS, javaMainClass);
084        }
085    }
086
087    public static void setupLauncherURIHandlerConf(Configuration launcherConf) {
088        for(Map.Entry<String, String> entry : Services.get().get(URIHandlerService.class).getLauncherConfig()) {
089            launcherConf.set(entry.getKey(), entry.getValue());
090        }
091    }
092
093    public static void setupMainArguments(Configuration launcherConf, String[] args) {
094        launcherConf.setInt(LauncherMapper.CONF_OOZIE_ACTION_MAIN_ARG_COUNT, args.length);
095        for (int i = 0; i < args.length; i++) {
096            launcherConf.set(LauncherMapper.CONF_OOZIE_ACTION_MAIN_ARG_PREFIX + i, args[i]);
097        }
098    }
099
100    public static void setupMaxOutputData(Configuration launcherConf, int maxOutputData) {
101        launcherConf.setInt(LauncherMapper.CONF_OOZIE_ACTION_MAX_OUTPUT_DATA, maxOutputData);
102    }
103
104    /**
105     * Set the maximum value of stats data
106     *
107     * @param launcherConf the oozie launcher configuration
108     * @param maxStatsData the maximum allowed size of stats data
109     */
110    public static void setupMaxExternalStatsSize(Configuration launcherConf, int maxStatsData){
111        launcherConf.setInt(LauncherMapper.CONF_OOZIE_EXTERNAL_STATS_MAX_SIZE, maxStatsData);
112    }
113
114    /**
115     * Set the maximum number of globbed files/dirs
116     *
117     * @param launcherConf the oozie launcher configuration
118     * @param fsGlobMax the maximum number of files/dirs for FS operation
119     */
120    public static void setupMaxFSGlob(Configuration launcherConf, int fsGlobMax){
121        launcherConf.setInt(LauncherMapper.CONF_OOZIE_ACTION_FS_GLOB_MAX, fsGlobMax);
122    }
123
124    public static void setupLauncherInfo(JobConf launcherConf, String jobId, String actionId, Path actionDir,
125            String recoveryId, Configuration actionConf, String prepareXML) throws IOException, HadoopAccessorException {
126
127        launcherConf.setMapperClass(LauncherMapper.class);
128        launcherConf.setSpeculativeExecution(false);
129        launcherConf.setNumMapTasks(1);
130        launcherConf.setNumReduceTasks(0);
131
132        launcherConf.set(LauncherMapper.OOZIE_JOB_ID, jobId);
133        launcherConf.set(LauncherMapper.OOZIE_ACTION_ID, actionId);
134        launcherConf.set(LauncherMapper.OOZIE_ACTION_DIR_PATH, actionDir.toString());
135        launcherConf.set(LauncherMapper.OOZIE_ACTION_RECOVERY_ID, recoveryId);
136        launcherConf.set(LauncherMapper.ACTION_PREPARE_XML, prepareXML);
137
138        actionConf.set(LauncherMapper.OOZIE_JOB_ID, jobId);
139        actionConf.set(LauncherMapper.OOZIE_ACTION_ID, actionId);
140
141        if (Services.get().getConf().getBoolean("oozie.hadoop-2.0.2-alpha.workaround.for.distributed.cache", false)) {
142          List<String> purgedEntries = new ArrayList<String>();
143          Collection<String> entries = actionConf.getStringCollection("mapreduce.job.cache.files");
144          for (String entry : entries) {
145            if (entry.contains("#")) {
146              purgedEntries.add(entry);
147            }
148          }
149          actionConf.setStrings("mapreduce.job.cache.files", purgedEntries.toArray(new String[purgedEntries.size()]));
150          launcherConf.setBoolean("oozie.hadoop-2.0.2-alpha.workaround.for.distributed.cache", true);
151        }
152
153        FileSystem fs =
154          Services.get().get(HadoopAccessorService.class).createFileSystem(launcherConf.get("user.name"),
155                                                                           actionDir.toUri(), launcherConf);
156        fs.mkdirs(actionDir);
157
158        OutputStream os = fs.create(new Path(actionDir, LauncherMapper.ACTION_CONF_XML));
159        try {
160            actionConf.writeXml(os);
161        } finally {
162            IOUtils.closeSafely(os);
163        }
164
165        launcherConf.setInputFormat(launcherInputFormatClassLocator.locateOrGet());
166        launcherConf.setOutputFormat(OozieLauncherOutputFormat.class);
167        launcherConf.setOutputCommitter(OozieLauncherOutputCommitter.class);
168    }
169
170    public static void setupYarnRestartHandling(JobConf launcherJobConf, Configuration actionConf, String launcherTag,
171                                                long launcherTime)
172            throws NoSuchAlgorithmException {
173        launcherJobConf.setLong(LauncherMainHadoopUtils.OOZIE_JOB_LAUNCH_TIME, launcherTime);
174        // Tags are limited to 100 chars so we need to hash them to make sure (the actionId otherwise doesn't have a max length)
175        String tag = getTag(launcherTag);
176        // keeping the oozie.child.mapreduce.job.tags instead of mapreduce.job.tags to avoid killing launcher itself.
177        // mapreduce.job.tags should only go to child job launch by launcher.
178        actionConf.set(LauncherMainHadoopUtils.CHILD_MAPREDUCE_JOB_TAGS, tag);
179    }
180
181    public static String getTag(String launcherTag) throws NoSuchAlgorithmException {
182        MessageDigest digest = MessageDigest.getInstance("MD5");
183        digest.update(launcherTag.getBytes(), 0, launcherTag.length());
184        String md5 = "oozie-" + new BigInteger(1, digest.digest()).toString(16);
185        return md5;
186    }
187
188    public static boolean isMainDone(RunningJob runningJob) throws IOException {
189        return runningJob.isComplete();
190    }
191
192    public static boolean isMainSuccessful(RunningJob runningJob) throws IOException {
193        boolean succeeded = runningJob.isSuccessful();
194        if (succeeded) {
195            Counters counters = runningJob.getCounters();
196            if (counters != null) {
197                Counters.Group group = counters.getGroup(LauncherMapper.COUNTER_GROUP);
198                if (group != null) {
199                    succeeded = group.getCounter(LauncherMapper.COUNTER_LAUNCHER_ERROR) == 0;
200                }
201            }
202        }
203        return succeeded;
204    }
205
206    /**
207     * Determine whether action has external child jobs or not
208     * @param actionData
209     * @return true/false
210     * @throws IOException
211     */
212    public static boolean hasExternalChildJobs(Map<String, String> actionData) throws IOException {
213        return actionData.containsKey(LauncherMapper.ACTION_DATA_EXTERNAL_CHILD_IDS);
214    }
215
216    /**
217     * Determine whether action has output data or not
218     * @param actionData
219     * @return true/false
220     * @throws IOException
221     */
222    public static boolean hasOutputData(Map<String, String> actionData) throws IOException {
223        return actionData.containsKey(LauncherMapper.ACTION_DATA_OUTPUT_PROPS);
224    }
225
226    /**
227     * Determine whether action has external stats or not
228     * @param actionData
229     * @return true/false
230     * @throws IOException
231     */
232    public static boolean hasStatsData(Map<String, String> actionData) throws IOException{
233        return actionData.containsKey(LauncherMapper.ACTION_DATA_STATS);
234    }
235
236    /**
237     * Determine whether action has new id (id swap) or not
238     * @param actionData
239     * @return true/false
240     * @throws IOException
241     */
242    public static boolean hasIdSwap(Map<String, String> actionData) throws IOException {
243        return actionData.containsKey(LauncherMapper.ACTION_DATA_NEW_ID);
244    }
245
246    /**
247     * Get the sequence file path storing all action data
248     * @param actionDir
249     * @return
250     */
251    public static Path getActionDataSequenceFilePath(Path actionDir) {
252        return new Path(actionDir, LauncherMapper.ACTION_DATA_SEQUENCE_FILE);
253    }
254
255    /**
256     * Utility function to load the contents of action data sequence file into
257     * memory object
258     *
259     * @param fs Action Filesystem
260     * @param actionDir Path
261     * @param conf Configuration
262     * @return Map action data
263     * @throws IOException
264     * @throws InterruptedException
265     */
266    public static Map<String, String> getActionData(final FileSystem fs, final Path actionDir, final Configuration conf)
267            throws IOException, InterruptedException {
268        UserGroupInformationService ugiService = Services.get().get(UserGroupInformationService.class);
269        UserGroupInformation ugi = ugiService.getProxyUser(conf.get(OozieClient.USER_NAME));
270
271        return ugi.doAs(new PrivilegedExceptionAction<Map<String, String>>() {
272            @Override
273            public Map<String, String> run() throws IOException {
274                Map<String, String> ret = new HashMap<String, String>();
275                Path seqFilePath = getActionDataSequenceFilePath(actionDir);
276                if (fs.exists(seqFilePath)) {
277                    SequenceFile.Reader seqFile = new SequenceFile.Reader(fs, seqFilePath, conf);
278                    Text key = new Text(), value = new Text();
279                    while (seqFile.next(key, value)) {
280                        ret.put(key.toString(), value.toString());
281                    }
282                    seqFile.close();
283                }
284                else { // maintain backward-compatibility. to be deprecated
285                    org.apache.hadoop.fs.FileStatus[] files = fs.listStatus(actionDir);
286                    InputStream is;
287                    BufferedReader reader = null;
288                    Properties props;
289                    if (files != null && files.length > 0) {
290                        for (int x = 0; x < files.length; x++) {
291                            Path file = files[x].getPath();
292                            if (file.equals(new Path(actionDir, "externalChildIds.properties"))) {
293                                is = fs.open(file);
294                                reader = new BufferedReader(new InputStreamReader(is));
295                                ret.put(LauncherMapper.ACTION_DATA_EXTERNAL_CHILD_IDS,
296                                        IOUtils.getReaderAsString(reader, -1));
297                            }
298                            else if (file.equals(new Path(actionDir, "newId.properties"))) {
299                                is = fs.open(file);
300                                reader = new BufferedReader(new InputStreamReader(is));
301                                props = PropertiesUtils.readProperties(reader, -1);
302                                ret.put(LauncherMapper.ACTION_DATA_NEW_ID, props.getProperty("id"));
303                            }
304                            else if (file.equals(new Path(actionDir, LauncherMapper.ACTION_DATA_OUTPUT_PROPS))) {
305                                int maxOutputData = conf.getInt(LauncherMapper.CONF_OOZIE_ACTION_MAX_OUTPUT_DATA,
306                                        2 * 1024);
307                                is = fs.open(file);
308                                reader = new BufferedReader(new InputStreamReader(is));
309                                ret.put(LauncherMapper.ACTION_DATA_OUTPUT_PROPS, PropertiesUtils
310                                        .propertiesToString(PropertiesUtils.readProperties(reader, maxOutputData)));
311                            }
312                            else if (file.equals(new Path(actionDir, LauncherMapper.ACTION_DATA_STATS))) {
313                                int statsMaxOutputData = conf.getInt(LauncherMapper.CONF_OOZIE_EXTERNAL_STATS_MAX_SIZE,
314                                        Integer.MAX_VALUE);
315                                is = fs.open(file);
316                                reader = new BufferedReader(new InputStreamReader(is));
317                                ret.put(LauncherMapper.ACTION_DATA_STATS, PropertiesUtils
318                                        .propertiesToString(PropertiesUtils.readProperties(reader, statsMaxOutputData)));
319                            }
320                            else if (file.equals(new Path(actionDir, LauncherMapper.ACTION_DATA_ERROR_PROPS))) {
321                                is = fs.open(file);
322                                reader = new BufferedReader(new InputStreamReader(is));
323                                ret.put(LauncherMapper.ACTION_DATA_ERROR_PROPS, IOUtils.getReaderAsString(reader, -1));
324                            }
325                        }
326                    }
327                }
328                return ret;
329            }
330        });
331    }
332
333    public static String getActionYarnTag(Configuration conf, String parentId, WorkflowAction wfAction) {
334        String tag;
335        if ( conf != null && conf.get(OOZIE_ACTION_YARN_TAG) != null) {
336            tag = conf.get(OOZIE_ACTION_YARN_TAG) + "@" + wfAction.getName();
337        } else if (parentId != null) {
338            tag = parentId + "@" + wfAction.getName();
339        } else {
340            tag = wfAction.getId();
341        }
342        return tag;
343    }
344
345}