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}