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.service; 020 021import org.apache.hadoop.io.Text; 022import org.apache.hadoop.mapred.JobClient; 023import org.apache.hadoop.mapred.JobConf; 024import org.apache.hadoop.fs.FileSystem; 025import org.apache.hadoop.fs.Path; 026import org.apache.hadoop.conf.Configuration; 027import org.apache.hadoop.mapreduce.security.token.delegation.DelegationTokenIdentifier; 028import org.apache.hadoop.net.NetUtils; 029import org.apache.hadoop.security.SecurityUtil; 030import org.apache.hadoop.security.UserGroupInformation; 031import org.apache.hadoop.security.token.Token; 032import org.apache.oozie.ErrorCode; 033import org.apache.oozie.action.hadoop.JavaActionExecutor; 034import org.apache.oozie.util.ParamChecker; 035import org.apache.oozie.util.XConfiguration; 036import org.apache.oozie.util.XLog; 037import org.apache.oozie.util.JobUtils; 038import org.apache.oozie.workflow.lite.LiteWorkflowAppParser; 039 040import java.io.File; 041import java.io.FileInputStream; 042import java.io.FilenameFilter; 043import java.io.IOException; 044import java.io.InputStream; 045import java.lang.reflect.InvocationTargetException; 046import java.lang.reflect.Method; 047import java.net.InetAddress; 048import java.net.URI; 049import java.net.URISyntaxException; 050import java.security.PrivilegedExceptionAction; 051import java.util.Arrays; 052import java.util.Comparator; 053import java.util.HashMap; 054import java.util.Map; 055import java.util.Set; 056import java.util.HashSet; 057import java.util.concurrent.ConcurrentHashMap; 058 059 060/** 061 * The HadoopAccessorService returns HadoopAccessor instances configured to work on behalf of a user-group. <p/> The 062 * default accessor used is the base accessor which just injects the UGI into the configuration instance used to 063 * create/obtain JobClient and FileSystem instances. 064 */ 065public class HadoopAccessorService implements Service { 066 067 private static XLog LOG = XLog.getLog(HadoopAccessorService.class); 068 069 public static final String CONF_PREFIX = Service.CONF_PREFIX + "HadoopAccessorService."; 070 public static final String JOB_TRACKER_WHITELIST = CONF_PREFIX + "jobTracker.whitelist"; 071 public static final String NAME_NODE_WHITELIST = CONF_PREFIX + "nameNode.whitelist"; 072 public static final String HADOOP_CONFS = CONF_PREFIX + "hadoop.configurations"; 073 public static final String ACTION_CONFS = CONF_PREFIX + "action.configurations"; 074 public static final String KERBEROS_AUTH_ENABLED = CONF_PREFIX + "kerberos.enabled"; 075 public static final String KERBEROS_KEYTAB = CONF_PREFIX + "keytab.file"; 076 public static final String KERBEROS_PRINCIPAL = CONF_PREFIX + "kerberos.principal"; 077 public static final Text MR_TOKEN_ALIAS = new Text("oozie mr token"); 078 079 protected static final String OOZIE_HADOOP_ACCESSOR_SERVICE_CREATED = "oozie.HadoopAccessorService.created"; 080 /** The Kerberos principal for the job tracker.*/ 081 protected static final String JT_PRINCIPAL = "mapreduce.jobtracker.kerberos.principal"; 082 /** The Kerberos principal for the resource manager.*/ 083 protected static final String RM_PRINCIPAL = "yarn.resourcemanager.principal"; 084 protected static final String HADOOP_JOB_TRACKER = "mapred.job.tracker"; 085 protected static final String HADOOP_JOB_TRACKER_2 = "mapreduce.jobtracker.address"; 086 protected static final String HADOOP_YARN_RM = "yarn.resourcemanager.address"; 087 private static final Map<String, Text> mrTokenRenewers = new HashMap<String, Text>(); 088 089 private static Configuration cachedConf; 090 091 private static final String DEFAULT_ACTIONNAME = "default"; 092 093 private Set<String> jobTrackerWhitelist = new HashSet<String>(); 094 private Set<String> nameNodeWhitelist = new HashSet<String>(); 095 private Map<String, Configuration> hadoopConfigs = new HashMap<String, Configuration>(); 096 private Map<String, File> actionConfigDirs = new HashMap<String, File>(); 097 private Map<String, Map<String, XConfiguration>> actionConfigs = new HashMap<String, Map<String, XConfiguration>>(); 098 099 private UserGroupInformationService ugiService; 100 101 /** 102 * Supported filesystem schemes for namespace federation 103 */ 104 public static final String SUPPORTED_FILESYSTEMS = CONF_PREFIX + "supported.filesystems"; 105 private Set<String> supportedSchemes; 106 private boolean allSchemesSupported; 107 108 public void init(Services services) throws ServiceException { 109 this.ugiService = services.get(UserGroupInformationService.class); 110 init(services.getConf()); 111 } 112 113 //for testing purposes, see XFsTestCase 114 public void init(Configuration conf) throws ServiceException { 115 for (String name : ConfigurationService.getStrings(conf, JOB_TRACKER_WHITELIST)) { 116 String tmp = name.toLowerCase().trim(); 117 if (tmp.length() == 0) { 118 continue; 119 } 120 jobTrackerWhitelist.add(tmp); 121 } 122 LOG.info( 123 "JOB_TRACKER_WHITELIST :" + jobTrackerWhitelist.toString() 124 + ", Total entries :" + jobTrackerWhitelist.size()); 125 for (String name : ConfigurationService.getStrings(conf, NAME_NODE_WHITELIST)) { 126 String tmp = name.toLowerCase().trim(); 127 if (tmp.length() == 0) { 128 continue; 129 } 130 nameNodeWhitelist.add(tmp); 131 } 132 LOG.info( 133 "NAME_NODE_WHITELIST :" + nameNodeWhitelist.toString() 134 + ", Total entries :" + nameNodeWhitelist.size()); 135 136 boolean kerberosAuthOn = ConfigurationService.getBoolean(conf, KERBEROS_AUTH_ENABLED); 137 LOG.info("Oozie Kerberos Authentication [{0}]", (kerberosAuthOn) ? "enabled" : "disabled"); 138 if (kerberosAuthOn) { 139 kerberosInit(conf); 140 } 141 else { 142 Configuration ugiConf = new Configuration(); 143 ugiConf.set("hadoop.security.authentication", "simple"); 144 UserGroupInformation.setConfiguration(ugiConf); 145 } 146 147 if (ugiService == null) { //for testing purposes, see XFsTestCase 148 this.ugiService = new UserGroupInformationService(); 149 } 150 151 loadHadoopConfigs(conf); 152 preLoadActionConfigs(conf); 153 154 supportedSchemes = new HashSet<String>(); 155 String[] schemesFromConf = ConfigurationService.getStrings(conf, SUPPORTED_FILESYSTEMS); 156 if(schemesFromConf != null) { 157 for (String scheme: schemesFromConf) { 158 scheme = scheme.trim(); 159 // If user gives "*", supportedSchemes will be empty, so that checking is not done i.e. all schemes allowed 160 if(scheme.equals("*")) { 161 if(schemesFromConf.length > 1) { 162 throw new ServiceException(ErrorCode.E0100, getClass().getName(), 163 SUPPORTED_FILESYSTEMS + " should contain either only wildcard or explicit list, not both"); 164 } 165 allSchemesSupported = true; 166 } 167 supportedSchemes.add(scheme); 168 } 169 } 170 171 setConfigForHadoopSecurityUtil(conf); 172 } 173 174 private void setConfigForHadoopSecurityUtil(Configuration conf) { 175 // Prior to HADOOP-12954 (2.9.0+), Hadoop sets hadoop.security.token.service.use_ip on startup in a static block with no 176 // way for Oozie to change it because Oozie doesn't load *-site.xml files on the classpath. HADOOP-12954 added a way to 177 // set this property via a setConfiguration method. Ideally, this would be part of JobClient so Oozie wouldn't have to 178 // worry about it and we could have different values for different clusters, but we can't; so we have to use the same value 179 // for every cluster Oozie is configured for. To that end, we'll use the default NN's configs. If that's not defined, 180 // we'll use the wildcard's configs. And if that's not defined, we'll use an arbitrary cluster's configs. In any case, 181 // if the version of Hadoop we're using doesn't include HADOOP-12954, we'll do nothing (there's no workaround), and 182 // hadoop.security.token.service.use_ip will have the default value. 183 String nameNode = conf.get(LiteWorkflowAppParser.DEFAULT_NAME_NODE); 184 if (nameNode != null) { 185 nameNode = nameNode.trim(); 186 if (nameNode.isEmpty()) { 187 nameNode = null; 188 } 189 } 190 if (nameNode == null && hadoopConfigs.containsKey("*")) { 191 nameNode = "*"; 192 } 193 if (nameNode == null) { 194 for (String nn : hadoopConfigs.keySet()) { 195 nn = nn.trim(); 196 if (!nn.isEmpty()) { 197 nameNode = nn; 198 break; 199 } 200 } 201 } 202 if (nameNode != null) { 203 Configuration hConf = getConfiguration(nameNode); 204 try { 205 Method setConfigurationMethod = SecurityUtil.class.getMethod("setConfiguration", Configuration.class); 206 setConfigurationMethod.invoke(null, hConf); 207 LOG.debug("Setting Hadoop SecurityUtil Configuration to that of {0}", nameNode); 208 } catch (NoSuchMethodException e) { 209 LOG.debug("Not setting Hadoop SecurityUtil Configuration because this version of Hadoop doesn't support it"); 210 } catch (Exception e) { 211 LOG.error("An Exception occurred while trying to call setConfiguration on {0} via Reflection. It won't be called.", 212 SecurityUtil.class.getName(), e); 213 } 214 } 215 } 216 217 private void kerberosInit(Configuration serviceConf) throws ServiceException { 218 try { 219 String keytabFile = ConfigurationService.get(serviceConf, KERBEROS_KEYTAB).trim(); 220 if (keytabFile.length() == 0) { 221 throw new ServiceException(ErrorCode.E0026, KERBEROS_KEYTAB); 222 } 223 String principal = SecurityUtil.getServerPrincipal( 224 serviceConf.get(KERBEROS_PRINCIPAL, "oozie/localhost@LOCALHOST"), 225 InetAddress.getLocalHost().getCanonicalHostName()); 226 if (principal.length() == 0) { 227 throw new ServiceException(ErrorCode.E0026, KERBEROS_PRINCIPAL); 228 } 229 Configuration conf = new Configuration(); 230 conf.set("hadoop.security.authentication", "kerberos"); 231 UserGroupInformation.setConfiguration(conf); 232 UserGroupInformation.loginUserFromKeytab(principal, keytabFile); 233 LOG.info("Got Kerberos ticket, keytab [{0}], Oozie principal principal [{1}]", 234 keytabFile, principal); 235 } 236 catch (ServiceException ex) { 237 throw ex; 238 } 239 catch (Exception ex) { 240 throw new ServiceException(ErrorCode.E0100, getClass().getName(), ex.getMessage(), ex); 241 } 242 } 243 244 private static final String[] HADOOP_CONF_FILES = 245 {"core-site.xml", "hdfs-site.xml", "mapred-site.xml", "yarn-site.xml", "hadoop-site.xml", "ssl-client.xml"}; 246 247 248 private Configuration loadHadoopConf(File dir) throws IOException { 249 Configuration hadoopConf = new XConfiguration(); 250 for (String file : HADOOP_CONF_FILES) { 251 File f = new File(dir, file); 252 if (f.exists()) { 253 InputStream is = new FileInputStream(f); 254 Configuration conf = new XConfiguration(is); 255 is.close(); 256 XConfiguration.copy(conf, hadoopConf); 257 } 258 } 259 return hadoopConf; 260 } 261 262 private Map<String, File> parseConfigDirs(String[] confDefs, String type) throws ServiceException, IOException { 263 Map<String, File> map = new HashMap<String, File>(); 264 File configDir = new File(ConfigurationService.getConfigurationDirectory()); 265 for (String confDef : confDefs) { 266 if (confDef.trim().length() > 0) { 267 String[] parts = confDef.split("="); 268 if (parts.length == 2) { 269 String hostPort = parts[0]; 270 String confDir = parts[1]; 271 File dir = new File(confDir); 272 if (!dir.isAbsolute()) { 273 dir = new File(configDir, confDir); 274 } 275 if (dir.exists()) { 276 map.put(hostPort.toLowerCase(), dir); 277 } 278 else { 279 throw new ServiceException(ErrorCode.E0100, getClass().getName(), 280 "could not find " + type + " configuration directory: " + 281 dir.getAbsolutePath()); 282 } 283 } 284 else { 285 throw new ServiceException(ErrorCode.E0100, getClass().getName(), 286 "Incorrect " + type + " configuration definition: " + confDef); 287 } 288 } 289 } 290 return map; 291 } 292 293 private void loadHadoopConfigs(Configuration serviceConf) throws ServiceException { 294 try { 295 Map<String, File> map = parseConfigDirs(ConfigurationService.getStrings(serviceConf, HADOOP_CONFS), 296 "hadoop"); 297 for (Map.Entry<String, File> entry : map.entrySet()) { 298 hadoopConfigs.put(entry.getKey(), loadHadoopConf(entry.getValue())); 299 } 300 } 301 catch (ServiceException ex) { 302 throw ex; 303 } 304 catch (Exception ex) { 305 throw new ServiceException(ErrorCode.E0100, getClass().getName(), ex.getMessage(), ex); 306 } 307 } 308 309 private void preLoadActionConfigs(Configuration serviceConf) throws ServiceException { 310 try { 311 actionConfigDirs = parseConfigDirs(ConfigurationService.getStrings(serviceConf, ACTION_CONFS), "action"); 312 for (String hostport : actionConfigDirs.keySet()) { 313 actionConfigs.put(hostport, new ConcurrentHashMap<String, XConfiguration>()); 314 } 315 } 316 catch (ServiceException ex) { 317 throw ex; 318 } 319 catch (Exception ex) { 320 throw new ServiceException(ErrorCode.E0100, getClass().getName(), ex.getMessage(), ex); 321 } 322 } 323 324 public void destroy() { 325 } 326 327 public Class<? extends Service> getInterface() { 328 return HadoopAccessorService.class; 329 } 330 331 private UserGroupInformation getUGI(String user) throws IOException { 332 return ugiService.getProxyUser(user); 333 } 334 335 /** 336 * Creates a JobConf using the site configuration for the specified hostname:port. 337 * <p/> 338 * If the specified hostname:port is not defined it falls back to the '*' site 339 * configuration if available. If the '*' site configuration is not available, 340 * the JobConf has all Hadoop defaults. 341 * 342 * @param hostPort hostname:port to lookup Hadoop site configuration. 343 * @return a JobConf with the corresponding site configuration for hostPort. 344 */ 345 public JobConf createJobConf(String hostPort) { 346 JobConf jobConf = new JobConf(getCachedConf()); 347 XConfiguration.copy(getConfiguration(hostPort), jobConf); 348 jobConf.setBoolean(OOZIE_HADOOP_ACCESSOR_SERVICE_CREATED, true); 349 return jobConf; 350 } 351 352 public Configuration getCachedConf() { 353 if (cachedConf == null) { 354 loadCachedConf(); 355 } 356 return cachedConf; 357 } 358 359 private void loadCachedConf() { 360 cachedConf = new Configuration(); 361 //for lazy loading 362 cachedConf.size(); 363 } 364 365 private XConfiguration loadActionConf(String hostPort, String action) { 366 File dir = actionConfigDirs.get(hostPort); 367 XConfiguration actionConf = new XConfiguration(); 368 if (dir != null) { 369 // See if a dir with the action name exists. If so, load all the xml files in the dir 370 File actionConfDir = new File(dir, action); 371 372 if (actionConfDir.exists() && actionConfDir.isDirectory()) { 373 LOG.info("Processing configuration files under [{0}]" 374 + " for action [{1}] and hostPort [{2}]", 375 actionConfDir.getAbsolutePath(), action, hostPort); 376 File[] xmlFiles = actionConfDir.listFiles( 377 new FilenameFilter() { 378 @Override 379 public boolean accept(File dir, String name) { 380 return name.endsWith(".xml"); 381 }}); 382 Arrays.sort(xmlFiles, new Comparator<File>() { 383 @Override 384 public int compare(File o1, File o2) { 385 return o1.getName().compareTo(o2.getName()); 386 } 387 }); 388 for (File f : xmlFiles) { 389 if (f.isFile() && f.canRead()) { 390 LOG.info("Processing configuration file [{0}]", f.getName()); 391 FileInputStream fis = null; 392 try { 393 fis = new FileInputStream(f); 394 XConfiguration conf = new XConfiguration(fis); 395 XConfiguration.copy(conf, actionConf); 396 } 397 catch (IOException ex) { 398 LOG 399 .warn("Could not read file [{0}] for action [{1}] configuration and hostPort [{2}]", 400 f.getAbsolutePath(), action, hostPort); 401 } 402 finally { 403 if (fis != null) { 404 try { fis.close(); } catch(IOException ioe) { } 405 } 406 } 407 } 408 } 409 } 410 } 411 412 // Now check for <action.xml> This way <action.xml> has priority over <action-dir>/*.xml 413 414 File actionConfFile = new File(dir, action + ".xml"); 415 if (actionConfFile.exists()) { 416 try { 417 XConfiguration conf = new XConfiguration(new FileInputStream(actionConfFile)); 418 XConfiguration.copy(conf, actionConf); 419 } 420 catch (IOException ex) { 421 LOG.warn("Could not read file [{0}] for action [{1}] configuration for hostPort [{2}]", 422 actionConfFile.getAbsolutePath(), action, hostPort); 423 } 424 } 425 426 return actionConf; 427 } 428 429 /** 430 * Returns a Configuration containing any defaults for an action for a particular cluster. 431 * <p/> 432 * This configuration is used as default for the action configuration and enables cluster 433 * level default values per action. 434 * 435 * @param hostPort hostname"port to lookup the action default confiugration. 436 * @param action action name. 437 * @return the default configuration for the action for the specified cluster. 438 */ 439 public XConfiguration createActionDefaultConf(String hostPort, String action) { 440 hostPort = (hostPort != null) ? hostPort.toLowerCase() : null; 441 Map<String, XConfiguration> hostPortActionConfigs = actionConfigs.get(hostPort); 442 if (hostPortActionConfigs == null) { 443 hostPortActionConfigs = actionConfigs.get("*"); 444 hostPort = "*"; 445 } 446 XConfiguration actionConf = hostPortActionConfigs.get(action); 447 if (actionConf == null) { 448 // doing lazy loading as we don't know upfront all actions, no need to synchronize 449 // as it is a read operation an in case of a race condition loading and inserting 450 // into the Map is idempotent and the action-config Map is a ConcurrentHashMap 451 452 // We first load a action of type default 453 // This allows for global configuration for all actions - for example 454 // all launchers in one queue and actions in another queue 455 // Are some configuration that applies to multiple actions - like 456 // config libraries path etc 457 actionConf = loadActionConf(hostPort, DEFAULT_ACTIONNAME); 458 459 // Action specific default configuration will override the default action config 460 461 XConfiguration.copy(loadActionConf(hostPort, action), actionConf); 462 hostPortActionConfigs.put(action, actionConf); 463 } 464 return new XConfiguration(actionConf.toProperties()); 465 } 466 467 private Configuration getConfiguration(String hostPort) { 468 hostPort = (hostPort != null) ? hostPort.toLowerCase() : null; 469 Configuration conf = hadoopConfigs.get(hostPort); 470 if (conf == null) { 471 conf = hadoopConfigs.get("*"); 472 if (conf == null) { 473 conf = new XConfiguration(); 474 } 475 } 476 return conf; 477 } 478 479 /** 480 * Return a JobClient created with the provided user/group. 481 * 482 * 483 * @param conf JobConf with all necessary information to create the 484 * JobClient. 485 * @return JobClient created with the provided user/group. 486 * @throws HadoopAccessorException if the client could not be created. 487 */ 488 public JobClient createJobClient(String user, final JobConf conf) throws HadoopAccessorException { 489 ParamChecker.notEmpty(user, "user"); 490 if (!conf.getBoolean(OOZIE_HADOOP_ACCESSOR_SERVICE_CREATED, false)) { 491 throw new HadoopAccessorException(ErrorCode.E0903); 492 } 493 String jobTracker = conf.get(JavaActionExecutor.HADOOP_JOB_TRACKER); 494 validateJobTracker(jobTracker); 495 try { 496 UserGroupInformation ugi = getUGI(user); 497 JobClient jobClient = ugi.doAs(new PrivilegedExceptionAction<JobClient>() { 498 public JobClient run() throws Exception { 499 return new JobClient(conf); 500 } 501 }); 502 Token<DelegationTokenIdentifier> mrdt = jobClient.getDelegationToken(getMRDelegationTokenRenewer(conf)); 503 conf.getCredentials().addToken(MR_TOKEN_ALIAS, mrdt); 504 return jobClient; 505 } 506 catch (InterruptedException ex) { 507 throw new HadoopAccessorException(ErrorCode.E0902, ex.getMessage(), ex); 508 } 509 catch (IOException ex) { 510 throw new HadoopAccessorException(ErrorCode.E0902, ex.getMessage(), ex); 511 } 512 } 513 514 /** 515 * Return a FileSystem created with the provided user for the specified URI. 516 * 517 * 518 * @param uri file system URI. 519 * @param conf Configuration with all necessary information to create the FileSystem. 520 * @return FileSystem created with the provided user/group. 521 * @throws HadoopAccessorException if the filesystem could not be created. 522 */ 523 public FileSystem createFileSystem(String user, final URI uri, final Configuration conf) 524 throws HadoopAccessorException { 525 ParamChecker.notEmpty(user, "user"); 526 if (!conf.getBoolean(OOZIE_HADOOP_ACCESSOR_SERVICE_CREATED, false)) { 527 throw new HadoopAccessorException(ErrorCode.E0903); 528 } 529 530 checkSupportedFilesystem(uri); 531 532 String nameNode = uri.getAuthority(); 533 if (nameNode == null) { 534 nameNode = conf.get("fs.default.name"); 535 if (nameNode != null) { 536 try { 537 nameNode = new URI(nameNode).getAuthority(); 538 } 539 catch (URISyntaxException ex) { 540 throw new HadoopAccessorException(ErrorCode.E0902, ex.getMessage(), ex); 541 } 542 } 543 } 544 validateNameNode(nameNode); 545 546 try { 547 UserGroupInformation ugi = getUGI(user); 548 return ugi.doAs(new PrivilegedExceptionAction<FileSystem>() { 549 public FileSystem run() throws Exception { 550 return FileSystem.get(uri, conf); 551 } 552 }); 553 } 554 catch (InterruptedException ex) { 555 throw new HadoopAccessorException(ErrorCode.E0902, ex.getMessage(), ex); 556 } 557 catch (IOException ex) { 558 throw new HadoopAccessorException(ErrorCode.E0902, ex.getMessage(), ex); 559 } 560 } 561 562 /** 563 * Validate Job tracker 564 * @param jobTrackerUri 565 * @throws HadoopAccessorException 566 */ 567 protected void validateJobTracker(String jobTrackerUri) throws HadoopAccessorException { 568 validate(jobTrackerUri, jobTrackerWhitelist, ErrorCode.E0900); 569 } 570 571 /** 572 * Validate Namenode list 573 * @param nameNodeUri 574 * @throws HadoopAccessorException 575 */ 576 protected void validateNameNode(String nameNodeUri) throws HadoopAccessorException { 577 validate(nameNodeUri, nameNodeWhitelist, ErrorCode.E0901); 578 } 579 580 private void validate(String uri, Set<String> whitelist, ErrorCode error) throws HadoopAccessorException { 581 if (uri != null) { 582 uri = uri.toLowerCase().trim(); 583 if (whitelist.size() > 0 && !whitelist.contains(uri)) { 584 throw new HadoopAccessorException(error, uri, whitelist); 585 } 586 } 587 } 588 589 public Text getMRDelegationTokenRenewer(JobConf jobConf) throws IOException { 590 if (UserGroupInformation.isSecurityEnabled()) { // secure cluster 591 return getMRTokenRenewerInternal(jobConf); 592 } 593 else { 594 return MR_TOKEN_ALIAS; //Doesn't matter what we pass as renewer 595 } 596 } 597 598 // Package private for unit test purposes 599 Text getMRTokenRenewerInternal(JobConf jobConf) throws IOException { 600 // Getting renewer correctly for JT principal also though JT in hadoop 1.x does not have 601 // support for renewing/cancelling tokens 602 String servicePrincipal = jobConf.get(RM_PRINCIPAL, jobConf.get(JT_PRINCIPAL)); 603 Text renewer; 604 if (servicePrincipal != null) { // secure cluster 605 renewer = mrTokenRenewers.get(servicePrincipal); 606 if (renewer == null) { 607 // Mimic org.apache.hadoop.mapred.Master.getMasterPrincipal() 608 String target = jobConf.get(HADOOP_YARN_RM, jobConf.get(HADOOP_JOB_TRACKER_2)); 609 if (target == null) { 610 target = jobConf.get(HADOOP_JOB_TRACKER); 611 } 612 try { 613 String addr = NetUtils.createSocketAddr(target).getHostName(); 614 renewer = new Text(SecurityUtil.getServerPrincipal(servicePrincipal, addr)); 615 LOG.info("Delegation Token Renewer details: Principal=" + servicePrincipal + ",Target=" + target 616 + ",Renewer=" + renewer); 617 } 618 catch (IllegalArgumentException iae) { 619 renewer = new Text(servicePrincipal.split("[/@]")[0]); 620 LOG.info("Delegation Token Renewer for " + servicePrincipal + " is " + renewer); 621 } 622 mrTokenRenewers.put(servicePrincipal, renewer); 623 } 624 } 625 else { 626 renewer = MR_TOKEN_ALIAS; //Doesn't matter what we pass as renewer 627 } 628 return renewer; 629 } 630 631 public void addFileToClassPath(String user, final Path file, final Configuration conf) 632 throws IOException { 633 ParamChecker.notEmpty(user, "user"); 634 try { 635 UserGroupInformation ugi = getUGI(user); 636 ugi.doAs(new PrivilegedExceptionAction<Void>() { 637 @Override 638 public Void run() throws Exception { 639 JobUtils.addFileToClassPath(file, conf, null); 640 return null; 641 } 642 }); 643 644 } 645 catch (InterruptedException ex) { 646 throw new IOException(ex); 647 } 648 649 } 650 651 /** 652 * checks configuration parameter if filesystem scheme is among the list of supported ones 653 * this makes system robust to filesystems other than HDFS also 654 */ 655 656 public void checkSupportedFilesystem(URI uri) throws HadoopAccessorException { 657 if (allSchemesSupported) 658 return; 659 String uriScheme = uri.getScheme(); 660 if (uriScheme != null) { // skip the check if no scheme is given 661 if(!supportedSchemes.isEmpty()) { 662 LOG.debug("Checking if filesystem " + uriScheme + " is supported"); 663 if (!supportedSchemes.contains(uriScheme)) { 664 throw new HadoopAccessorException(ErrorCode.E0904, uriScheme, uri.toString()); 665 } 666 } 667 } 668 } 669 670 public Set<String> getSupportedSchemes() { 671 return supportedSchemes; 672 } 673 674}