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 java.io.File; 022import java.io.FileInputStream; 023import java.io.IOException; 024import java.io.InputStream; 025import java.io.OutputStream; 026import java.net.URI; 027import java.net.URISyntaxException; 028import java.net.URL; 029import java.net.URLDecoder; 030import java.text.MessageFormat; 031import java.text.ParseException; 032import java.text.SimpleDateFormat; 033import java.util.ArrayList; 034import java.util.Arrays; 035import java.util.Calendar; 036import java.util.Comparator; 037import java.util.Date; 038import java.util.Enumeration; 039import java.util.HashMap; 040import java.util.HashSet; 041import java.util.List; 042import java.util.Map; 043import java.util.Properties; 044import java.util.Set; 045import java.util.TimeZone; 046import java.util.Map.Entry; 047import org.apache.commons.lang.StringUtils; 048import org.apache.hadoop.conf.Configuration; 049import org.apache.hadoop.fs.FileStatus; 050import org.apache.hadoop.fs.FileSystem; 051import org.apache.hadoop.fs.Path; 052import org.apache.hadoop.fs.PathFilter; 053import org.apache.hadoop.fs.permission.FsPermission; 054import org.apache.hadoop.io.IOUtils; 055import org.apache.oozie.action.ActionExecutor; 056import org.apache.oozie.action.hadoop.JavaActionExecutor; 057import org.apache.oozie.client.rest.JsonUtils; 058import org.apache.oozie.hadoop.utils.HadoopShims; 059import org.apache.oozie.util.Instrumentable; 060import org.apache.oozie.util.Instrumentation; 061import org.apache.oozie.util.XConfiguration; 062import org.apache.oozie.util.XLog; 063import com.google.common.annotations.VisibleForTesting; 064 065import org.apache.oozie.ErrorCode; 066import org.jdom.JDOMException; 067 068public class ShareLibService implements Service, Instrumentable { 069 070 public static final String LAUNCHERJAR_LIB_RETENTION = CONF_PREFIX + "ShareLibService.temp.sharelib.retention.days"; 071 072 public static final String SHARELIB_MAPPING_FILE = CONF_PREFIX + "ShareLibService.mapping.file"; 073 074 public static final String SHIP_LAUNCHER_JAR = "oozie.action.ship.launcher.jar"; 075 076 public static final String PURGE_INTERVAL = CONF_PREFIX + "ShareLibService.purge.interval"; 077 078 public static final String FAIL_FAST_ON_STARTUP = CONF_PREFIX + "ShareLibService.fail.fast.on.startup"; 079 080 private static final String PERMISSION_STRING = "-rwxr-xr-x"; 081 082 public static final String LAUNCHER_LIB_PREFIX = "launcher_"; 083 084 public static final String SHARE_LIB_PREFIX = "lib_"; 085 086 public static final SimpleDateFormat dateFormat = new SimpleDateFormat("yyyyMMddHHmmss"); 087 088 private Services services; 089 090 private Map<String, List<Path>> shareLibMap = new HashMap<String, List<Path>>(); 091 092 private Map<String, Map<Path, Configuration>> shareLibConfigMap = new HashMap<String, Map<Path, Configuration>>(); 093 094 private Map<String, List<Path>> launcherLibMap = new HashMap<String, List<Path>>(); 095 096 private Set<String> actionConfSet = new HashSet<String>(); 097 098 // symlink mapping. Oozie keeps on checking symlink path and if changes, Oozie reloads the sharelib 099 private Map<String, Map<Path, Path>> symlinkMapping = new HashMap<String, Map<Path, Path>>(); 100 101 private static XLog LOG = XLog.getLog(ShareLibService.class); 102 103 private String sharelibMappingFile; 104 105 private boolean isShipLauncherEnabled = false; 106 107 public static String SHARE_LIB_CONF_PREFIX = "oozie"; 108 109 private boolean shareLibLoadAttempted = false; 110 111 private String sharelibMetaFileOldTimeStamp; 112 113 private String sharelibDirOld; 114 115 FileSystem fs; 116 117 final long retentionTime = 1000 * 60 * 60 * 24 * ConfigurationService.getInt(LAUNCHERJAR_LIB_RETENTION); 118 119 @Override 120 public void init(Services services) throws ServiceException { 121 this.services = services; 122 sharelibMappingFile = ConfigurationService.get(services.getConf(), SHARELIB_MAPPING_FILE); 123 isShipLauncherEnabled = ConfigurationService.getBoolean(services.getConf(), SHIP_LAUNCHER_JAR); 124 boolean failOnfailure = ConfigurationService.getBoolean(services.getConf(), FAIL_FAST_ON_STARTUP); 125 Path launcherlibPath = getLauncherlibPath(); 126 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 127 URI uri = launcherlibPath.toUri(); 128 try { 129 fs = FileSystem.get(has.createJobConf(uri.getAuthority())); 130 //cache action key sharelib conf list 131 cacheActionKeySharelibConfList(); 132 updateLauncherLib(); 133 updateShareLib(); 134 } 135 catch (Throwable e) { 136 if (failOnfailure) { 137 LOG.error("Sharelib initialization fails", e); 138 throw new ServiceException(ErrorCode.E0104, getClass().getName(), "Sharelib initialization fails. ", e); 139 } 140 else { 141 // We don't want to actually fail init by throwing an Exception, so only create the ServiceException and 142 // log it 143 ServiceException se = new ServiceException(ErrorCode.E0104, getClass().getName(), 144 "Not able to cache sharelib. An Admin needs to install the sharelib with oozie-setup.sh and issue the " 145 + "'oozie admin' CLI command to update the sharelib", e); 146 LOG.error(se); 147 } 148 } 149 Runnable purgeLibsRunnable = new Runnable() { 150 @Override 151 public void run() { 152 System.out.flush(); 153 try { 154 // Only one server should purge sharelib 155 if (Services.get().get(JobsConcurrencyService.class).isLeader()) { 156 final Date current = Calendar.getInstance(TimeZone.getTimeZone("GMT")).getTime(); 157 purgeLibs(fs, LAUNCHER_LIB_PREFIX, current); 158 purgeLibs(fs, SHARE_LIB_PREFIX, current); 159 } 160 } 161 catch (IOException e) { 162 LOG.error("There was an issue purging the sharelib", e); 163 } 164 } 165 }; 166 services.get(SchedulerService.class).schedule(purgeLibsRunnable, 10, 167 ConfigurationService.getInt(services.getConf(), PURGE_INTERVAL) * 60 * 60 * 24, 168 SchedulerService.Unit.SEC); 169 } 170 171 /** 172 * Recursively change permissions. 173 * 174 * @throws IOException Signals that an I/O exception has occurred. 175 */ 176 private void updateLauncherLib() throws IOException { 177 if (isShipLauncherEnabled) { 178 if (fs == null) { 179 Path launcherlibPath = getLauncherlibPath(); 180 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 181 URI uri = launcherlibPath.toUri(); 182 fs = FileSystem.get(has.createJobConf(uri.getAuthority())); 183 } 184 Path launcherlibPath = getLauncherlibPath(); 185 setupLauncherLibPath(fs, launcherlibPath); 186 recursiveChangePermissions(fs, launcherlibPath, FsPermission.valueOf(PERMISSION_STRING)); 187 } 188 189 } 190 191 /** 192 * Copy launcher jars to Temp directory. 193 * 194 * @param fs the FileSystem 195 * @param tmpLauncherLibPath the tmp launcher lib path 196 * @throws IOException Signals that an I/O exception has occurred. 197 */ 198 private void setupLauncherLibPath(FileSystem fs, Path tmpLauncherLibPath) throws IOException { 199 200 ActionService actionService = Services.get().get(ActionService.class); 201 List<Class> classes = JavaActionExecutor.getCommonLauncherClasses(); 202 Path baseDir = new Path(tmpLauncherLibPath, JavaActionExecutor.OOZIE_COMMON_LIBDIR); 203 copyJarContainingClasses(classes, fs, baseDir, JavaActionExecutor.OOZIE_COMMON_LIBDIR); 204 Set<String> actionTypes = actionService.getActionTypes(); 205 for (String key : actionTypes) { 206 ActionExecutor executor = actionService.getExecutor(key); 207 if (executor instanceof JavaActionExecutor) { 208 JavaActionExecutor jexecutor = (JavaActionExecutor) executor; 209 classes = jexecutor.getLauncherClasses(); 210 if (classes != null) { 211 String type = executor.getType(); 212 Path executorDir = new Path(tmpLauncherLibPath, type); 213 copyJarContainingClasses(classes, fs, executorDir, type); 214 } 215 } 216 } 217 } 218 219 /** 220 * Recursive change permissions. 221 * 222 * @param fs the FileSystem 223 * @param path the Path 224 * @param perm is permission 225 * @throws IOException Signals that an I/O exception has occurred. 226 */ 227 private void recursiveChangePermissions(FileSystem fs, Path path, FsPermission fsPerm) throws IOException { 228 fs.setPermission(path, fsPerm); 229 FileStatus[] filesStatus = fs.listStatus(path); 230 for (int i = 0; i < filesStatus.length; i++) { 231 Path p = filesStatus[i].getPath(); 232 if (filesStatus[i].isDir()) { 233 recursiveChangePermissions(fs, p, fsPerm); 234 } 235 else { 236 fs.setPermission(p, fsPerm); 237 } 238 } 239 } 240 241 /** 242 * Copy jar containing classes. 243 * 244 * @param classes the classes 245 * @param fs the FileSystem 246 * @param executorDir is Path 247 * @param type is sharelib key 248 * @throws IOException Signals that an I/O exception has occurred. 249 */ 250 private void copyJarContainingClasses(List<Class> classes, FileSystem fs, Path executorDir, String type) 251 throws IOException { 252 fs.mkdirs(executorDir); 253 Set<String> localJarSet = new HashSet<String>(); 254 for (Class c : classes) { 255 String localJar = findContainingJar(c); 256 if (localJar != null) { 257 localJarSet.add(localJar); 258 } 259 else { 260 throw new IOException("No jar containing " + c + " found"); 261 } 262 } 263 List<Path> listOfPaths = new ArrayList<Path>(); 264 for (String localJarStr : localJarSet) { 265 File localJar = new File(localJarStr); 266 copyFromLocalFile(localJar, fs, executorDir); 267 Path path = new Path(executorDir, localJar.getName()); 268 listOfPaths.add(path); 269 LOG.info(localJar.getName() + " uploaded to " + executorDir.toString()); 270 } 271 launcherLibMap.put(type, listOfPaths); 272 273 } 274 275 private static boolean copyFromLocalFile(File src, FileSystem dstFS, Path dstDir) throws IOException { 276 Path dst = new Path(dstDir, src.getName()); 277 InputStream in=null; 278 OutputStream out = null; 279 try { 280 in = new FileInputStream(src); 281 out = dstFS.create(dst, true); 282 IOUtils.copyBytes(in, out, dstFS.getConf(), true); 283 } catch (IOException e) { 284 IOUtils.closeStream(out); 285 IOUtils.closeStream(in); 286 throw e; 287 } 288 return true; 289 290 } 291 292 /** 293 * Gets the path recursively. 294 * 295 * @param fs the FileSystem 296 * @param rootDir the root directory 297 * @param listOfPaths the list of paths 298 * @param shareLibKey the share lib key 299 * @return the path recursively 300 * @throws IOException Signals that an I/O exception has occurred. 301 */ 302 private void getPathRecursively(FileSystem fs, Path rootDir, List<Path> listOfPaths, String shareLibKey, 303 Map<String, Map<Path, Configuration>> shareLibConfigMap) throws IOException { 304 if (rootDir == null) { 305 return; 306 } 307 308 try { 309 if (fs.isFile(new Path(new URI(rootDir.toString()).getPath()))) { 310 Path filePath = new Path(new URI(rootDir.toString()).getPath()); 311 312 if (isFilePartOfConfList(rootDir)) { 313 cachePropertyFile(filePath, shareLibKey, shareLibConfigMap); 314 } 315 316 listOfPaths.add(rootDir); 317 return; 318 } 319 320 FileStatus[] status = fs.listStatus(rootDir); 321 if (status == null) { 322 LOG.info("Shared lib " + rootDir + " doesn't exist, not adding to cache"); 323 return; 324 } 325 326 for (FileStatus file : status) { 327 if (file.isDir()) { 328 getPathRecursively(fs, file.getPath(), listOfPaths, shareLibKey, shareLibConfigMap); 329 } 330 else { 331 if (isFilePartOfConfList(file.getPath())) { 332 cachePropertyFile(file.getPath(), shareLibKey, shareLibConfigMap); 333 } 334 listOfPaths.add(file.getPath()); 335 } 336 } 337 } 338 catch (URISyntaxException e) { 339 throw new IOException(e); 340 } 341 catch (JDOMException e) { 342 throw new IOException(e); 343 } 344 } 345 346 public Map<String, List<Path>> getShareLib() { 347 return shareLibMap; 348 } 349 350 private Map<String, Map<Path, Path>> getSymlinkMapping() { 351 return symlinkMapping; 352 } 353 354 /** 355 * Gets the action sharelib lib jars. 356 * 357 * @param shareLibKey the sharelib key 358 * @return List of paths 359 * @throws IOException Signals that an I/O exception has occurred. 360 */ 361 public List<Path> getShareLibJars(String shareLibKey) throws IOException { 362 // Sharelib map is empty means that on previous or startup attempt of 363 // caching sharelib has failed.Trying to reload 364 if (shareLibMap.isEmpty() && !shareLibLoadAttempted) { 365 synchronized (ShareLibService.class) { 366 if (shareLibMap.isEmpty()) { 367 updateShareLib(); 368 shareLibLoadAttempted = true; 369 } 370 } 371 } 372 checkSymlink(shareLibKey); 373 return shareLibMap.get(shareLibKey); 374 } 375 376 private void checkSymlink(String shareLibKey) throws IOException { 377 if (!HadoopShims.isSymlinkSupported() || symlinkMapping.get(shareLibKey) == null 378 || symlinkMapping.get(shareLibKey).isEmpty()) { 379 return; 380 } 381 382 HadoopShims fileSystem = new HadoopShims(fs); 383 for (Path path : symlinkMapping.get(shareLibKey).keySet()) { 384 if (!symlinkMapping.get(shareLibKey).get(path).equals(fileSystem.getSymLinkTarget(path))) { 385 synchronized (ShareLibService.class) { 386 Map<String, List<Path>> tmpShareLibMap = new HashMap<String, List<Path>>(shareLibMap); 387 388 Map<String, Map<Path, Configuration>> tmpShareLibConfigMap = new HashMap<String, Map<Path, Configuration>>( 389 shareLibConfigMap); 390 391 Map<String, Map<Path, Path>> tmpSymlinkMapping = new HashMap<String, Map<Path, Path>>( 392 symlinkMapping); 393 394 LOG.info(MessageFormat.format("Symlink target for [{0}] has changed, was [{1}], now [{2}]", 395 shareLibKey, path, fileSystem.getSymLinkTarget(path))); 396 loadShareLibMetaFile(tmpShareLibMap, tmpSymlinkMapping, tmpShareLibConfigMap, sharelibMappingFile, 397 shareLibKey); 398 shareLibMap = tmpShareLibMap; 399 symlinkMapping = tmpSymlinkMapping; 400 shareLibConfigMap = tmpShareLibConfigMap; 401 return; 402 } 403 404 } 405 } 406 407 } 408 409 /** 410 * Gets the launcher jars. 411 * 412 * @param shareLibKey the shareLib key 413 * @return launcher jars paths 414 * @throws IOException Signals that an I/O exception has occurred. 415 */ 416 public List<Path> getSystemLibJars(String shareLibKey) throws IOException { 417 List<Path> returnList = new ArrayList<Path>(); 418 // Sharelib map is empty means that on previous or startup attempt of 419 // caching launcher jars has failed.Trying to reload 420 if (isShipLauncherEnabled) { 421 if (launcherLibMap.isEmpty()) { 422 synchronized (ShareLibService.class) { 423 if (launcherLibMap.isEmpty()) { 424 updateLauncherLib(); 425 } 426 } 427 } 428 if (launcherLibMap.get(shareLibKey) != null) { 429 returnList.addAll(launcherLibMap.get(shareLibKey)); 430 } 431 } 432 if (shareLibKey.equals(JavaActionExecutor.OOZIE_COMMON_LIBDIR)) { 433 List<Path> sharelibList = getShareLibJars(shareLibKey); 434 if (sharelibList != null) { 435 returnList.addAll(sharelibList); 436 } 437 } 438 return returnList; 439 } 440 441 /** 442 * Find containing jar containing. 443 * 444 * @param clazz the clazz 445 * @return the string 446 */ 447 @VisibleForTesting 448 protected String findContainingJar(Class clazz) { 449 ClassLoader loader = clazz.getClassLoader(); 450 String classFile = clazz.getName().replaceAll("\\.", "/") + ".class"; 451 try { 452 for (Enumeration itr = loader.getResources(classFile); itr.hasMoreElements();) { 453 URL url = (URL) itr.nextElement(); 454 if ("jar".equals(url.getProtocol())) { 455 String toReturn = url.getPath(); 456 if (toReturn.startsWith("file:")) { 457 toReturn = toReturn.substring("file:".length()); 458 // URLDecoder is a misnamed class, since it actually 459 // decodes 460 // x-www-form-urlencoded MIME type rather than actual 461 // URL encoding (which the file path has). Therefore it 462 // would 463 // decode +s to ' 's which is incorrect (spaces are 464 // actually 465 // either unencoded or encoded as "%20"). Replace +s 466 // first, so 467 // that they are kept sacred during the decoding 468 // process. 469 toReturn = toReturn.replaceAll("\\+", "%2B"); 470 toReturn = URLDecoder.decode(toReturn, "UTF-8"); 471 toReturn = toReturn.replaceAll("!.*$", ""); 472 return toReturn; 473 } 474 } 475 } 476 } 477 catch (IOException ioe) { 478 throw new RuntimeException(ioe); 479 } 480 return null; 481 } 482 483 /** 484 * Purge libs. 485 * 486 * @param fs the fs 487 * @param prefix the prefix 488 * @param current the current time 489 * @throws IOException Signals that an I/O exception has occurred. 490 */ 491 private void purgeLibs(FileSystem fs, final String prefix, final Date current) throws IOException { 492 Path executorLibBasePath = services.get(WorkflowAppService.class).getSystemLibPath(); 493 PathFilter directoryFilter = new PathFilter() { 494 @Override 495 public boolean accept(Path path) { 496 if (path.getName().startsWith(prefix)) { 497 String name = path.getName(); 498 String time = name.substring(prefix.length()); 499 Date d = null; 500 try { 501 d = dateFormat.parse(time); 502 } 503 catch (ParseException e) { 504 return false; 505 } 506 return (current.getTime() - d.getTime()) > retentionTime; 507 } 508 else { 509 return false; 510 } 511 } 512 }; 513 FileStatus[] dirList = fs.listStatus(executorLibBasePath, directoryFilter); 514 Arrays.sort(dirList, new Comparator<FileStatus>() { 515 // sort in desc order 516 @Override 517 public int compare(FileStatus o1, FileStatus o2) { 518 return o2.getPath().getName().compareTo(o1.getPath().getName()); 519 } 520 }); 521 522 // Logic is to keep all share-lib between current timestamp and 7days old + 1 latest sharelib older than 7 days. 523 // refer OOZIE-1761 524 for (int i = 1; i < dirList.length; i++) { 525 Path dirPath = dirList[i].getPath(); 526 fs.delete(dirPath, true); 527 LOG.info("Deleted old launcher jar lib directory {0}", dirPath.getName()); 528 } 529 } 530 531 @Override 532 public void destroy() { 533 shareLibMap.clear(); 534 launcherLibMap.clear(); 535 } 536 537 @Override 538 public Class<? extends Service> getInterface() { 539 return ShareLibService.class; 540 } 541 542 /** 543 * Update share lib cache. 544 * 545 * @return the map 546 * @throws IOException Signals that an I/O exception has occurred. 547 */ 548 public Map<String, String> updateShareLib() throws IOException { 549 Map<String, String> status = new HashMap<String, String>(); 550 551 if (fs == null) { 552 Path launcherlibPath = getLauncherlibPath(); 553 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 554 URI uri = launcherlibPath.toUri(); 555 fs = FileSystem.get(has.createJobConf(uri.getAuthority())); 556 } 557 558 Map<String, List<Path>> tempShareLibMap = new HashMap<String, List<Path>>(); 559 Map<String, Map<Path, Path>> tmpSymlinkMapping = new HashMap<String, Map<Path, Path>>(); 560 Map<String, Map<Path, Configuration>> tmpShareLibConfigMap = new HashMap<String, Map<Path, Configuration>>(); 561 562 if (!StringUtils.isEmpty(sharelibMappingFile.trim())) { 563 String sharelibMetaFileNewTimeStamp = JsonUtils.formatDateRfc822( 564 new Date(fs.getFileStatus(new Path(sharelibMappingFile)).getModificationTime()), "GMT"); 565 loadShareLibMetaFile(tempShareLibMap, tmpSymlinkMapping, tmpShareLibConfigMap, sharelibMappingFile, null); 566 status.put("sharelibMetaFile", sharelibMappingFile); 567 status.put("sharelibMetaFileNewTimeStamp", sharelibMetaFileNewTimeStamp); 568 status.put("sharelibMetaFileOldTimeStamp", sharelibMetaFileOldTimeStamp); 569 sharelibMetaFileOldTimeStamp = sharelibMetaFileNewTimeStamp; 570 } 571 else { 572 Path shareLibpath = getLatestLibPath(services.get(WorkflowAppService.class).getSystemLibPath(), 573 SHARE_LIB_PREFIX); 574 loadShareLibfromDFS(tempShareLibMap, shareLibpath, tmpShareLibConfigMap); 575 576 if (shareLibpath != null) { 577 status.put("sharelibDirNew", shareLibpath.toString()); 578 status.put("sharelibDirOld", sharelibDirOld); 579 sharelibDirOld = shareLibpath.toString(); 580 } 581 582 } 583 shareLibMap = tempShareLibMap; 584 symlinkMapping = tmpSymlinkMapping; 585 shareLibConfigMap = tmpShareLibConfigMap; 586 return status; 587 } 588 589 /** 590 * Update share lib cache. Parse the share lib directory and each sub directory is a action key 591 * 592 * @param shareLibMap the share lib jar map 593 * @param shareLibpath the share libpath 594 * @throws IOException Signals that an I/O exception has occurred. 595 */ 596 private void loadShareLibfromDFS(Map<String, List<Path>> shareLibMap, Path shareLibpath, 597 Map<String, Map<Path, Configuration>> shareLibConfigMap) throws IOException { 598 599 if (shareLibpath == null) { 600 LOG.info("No share lib directory found"); 601 return; 602 603 } 604 605 FileStatus[] dirList = fs.listStatus(shareLibpath); 606 607 if (dirList == null) { 608 return; 609 } 610 611 for (FileStatus dir : dirList) { 612 if (!dir.isDir()) { 613 continue; 614 } 615 List<Path> listOfPaths = new ArrayList<Path>(); 616 getPathRecursively(fs, dir.getPath(), listOfPaths, dir.getPath().getName(), shareLibConfigMap); 617 shareLibMap.put(dir.getPath().getName(), listOfPaths); 618 LOG.info("Share lib for " + dir.getPath().getName() + ":" + listOfPaths); 619 620 } 621 622 } 623 624 /** 625 * Load share lib text file. Sharelib mapping files contains list of key=value. where key is the action key and 626 * value is the DFS location of sharelib files. 627 * 628 * @param shareLibMap the share lib jar map 629 * @param symlinkMapping the symlink mapping 630 * @param sharelibFileMapping the sharelib file mapping 631 * @param shareLibKey the share lib key 632 * @throws IOException Signals that an I/O exception has occurred. 633 * @parm shareLibKey the sharelib key 634 */ 635 private void loadShareLibMetaFile(Map<String, List<Path>> shareLibMap, Map<String, Map<Path, Path>> symlinkMapping, 636 Map<String, Map<Path, Configuration>> shareLibConfigMap, String sharelibFileMapping, String shareLibKey) 637 throws IOException { 638 639 Path shareFileMappingPath = new Path(sharelibFileMapping); 640 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 641 FileSystem filesystem = FileSystem.get(has.createJobConf(shareFileMappingPath.toUri().getAuthority())); 642 Properties prop = new Properties(); 643 prop.load(filesystem.open(new Path(sharelibFileMapping))); 644 645 for (Object keyObject : prop.keySet()) { 646 String key = (String) keyObject; 647 String mapKey = key.substring(SHARE_LIB_CONF_PREFIX.length() + 1); 648 if (key.toLowerCase().startsWith(SHARE_LIB_CONF_PREFIX) 649 && (shareLibKey == null || shareLibKey.equals(mapKey))) { 650 loadSharelib(shareLibMap, symlinkMapping, shareLibConfigMap, mapKey, 651 ((String) prop.get(key)).split(",")); 652 } 653 } 654 } 655 656 private void loadSharelib(Map<String, List<Path>> tmpShareLibMap, Map<String, Map<Path, Path>> tmpSymlinkMapping, 657 Map<String, Map<Path, Configuration>> shareLibConfigMap, String shareLibKey, String pathList[]) 658 throws IOException { 659 List<Path> listOfPaths = new ArrayList<Path>(); 660 Map<Path, Path> symlinkMappingforAction = new HashMap<Path, Path>(); 661 HadoopShims fileSystem = new HadoopShims(fs); 662 663 for (String dfsPath : pathList) { 664 Path path = new Path(dfsPath); 665 getPathRecursively(fs, new Path(dfsPath), listOfPaths, shareLibKey, shareLibConfigMap); 666 if (HadoopShims.isSymlinkSupported() && fileSystem.isSymlink(path)) { 667 symlinkMappingforAction.put(path, fileSystem.getSymLinkTarget(path)); 668 } 669 } 670 if (HadoopShims.isSymlinkSupported()) { 671 LOG.info("symlink for " + shareLibKey + ":" + symlinkMappingforAction); 672 tmpSymlinkMapping.put(shareLibKey, symlinkMappingforAction); 673 } 674 tmpShareLibMap.put(shareLibKey, listOfPaths); 675 LOG.info("Share lib for " + shareLibKey + ":" + listOfPaths); 676 } 677 678 /** 679 * Gets the launcherlib path. 680 * 681 * @return the launcherlib path 682 */ 683 private Path getLauncherlibPath() { 684 String formattedDate = dateFormat.format(Calendar.getInstance(TimeZone.getTimeZone("GMT")).getTime()); 685 Path tmpLauncherLibPath = new Path(services.get(WorkflowAppService.class).getSystemLibPath(), LAUNCHER_LIB_PREFIX 686 + formattedDate); 687 return tmpLauncherLibPath; 688 } 689 690 /** 691 * Gets the Latest lib path. 692 * 693 * @param rootDir the root dir 694 * @param prefix the prefix 695 * @return latest lib path 696 * @throws IOException Signals that an I/O exception has occurred. 697 */ 698 public Path getLatestLibPath(Path rootDir, final String prefix) throws IOException { 699 Date max = new Date(0L); 700 Path path = null; 701 PathFilter directoryFilter = new PathFilter() { 702 @Override 703 public boolean accept(Path path) { 704 return path.getName().startsWith(prefix); 705 } 706 }; 707 708 FileStatus[] files = fs.listStatus(rootDir, directoryFilter); 709 for (FileStatus file : files) { 710 String name = file.getPath().getName().toString(); 711 String time = name.substring(prefix.length()); 712 Date d = null; 713 try { 714 d = dateFormat.parse(time); 715 } 716 catch (ParseException e) { 717 continue; 718 } 719 if (d.compareTo(max) > 0) { 720 path = file.getPath(); 721 max = d; 722 } 723 } 724 // If there are no timestamped directories, fall back to root directory 725 if (path == null) { 726 path = rootDir; 727 } 728 return path; 729 } 730 731 /** 732 * Instruments the log service. 733 * <p/> 734 * It sets instrumentation variables indicating the location of the sharelib and launcherlib 735 * 736 * @param instr instrumentation to use. 737 */ 738 @Override 739 public void instrument(Instrumentation instr) { 740 instr.addVariable("libs", "sharelib.source", new Instrumentation.Variable<String>() { 741 @Override 742 public String getValue() { 743 if (!StringUtils.isEmpty(sharelibMappingFile.trim())) { 744 return SHARELIB_MAPPING_FILE; 745 } 746 return WorkflowAppService.SYSTEM_LIB_PATH; 747 } 748 }); 749 instr.addVariable("libs", "sharelib.mapping.file", new Instrumentation.Variable<String>() { 750 @Override 751 public String getValue() { 752 if (!StringUtils.isEmpty(sharelibMappingFile.trim())) { 753 return sharelibMappingFile; 754 } 755 return "(none)"; 756 } 757 }); 758 instr.addVariable("libs", "sharelib.system.libpath", new Instrumentation.Variable<String>() { 759 @Override 760 public String getValue() { 761 String sharelibPath = "(unavailable)"; 762 try { 763 Path libPath = getLatestLibPath(services.get(WorkflowAppService.class).getSystemLibPath(), 764 SHARE_LIB_PREFIX); 765 if (libPath != null) { 766 sharelibPath = libPath.toUri().toString(); 767 } 768 } 769 catch (IOException ioe) { 770 // ignore exception because we're just doing instrumentation 771 } 772 return sharelibPath; 773 } 774 }); 775 instr.addVariable("libs", "sharelib.mapping.file.timestamp", new Instrumentation.Variable<String>() { 776 @Override 777 public String getValue() { 778 if (!StringUtils.isEmpty(sharelibMetaFileOldTimeStamp)) { 779 return sharelibMetaFileOldTimeStamp; 780 } 781 return "(none)"; 782 } 783 }); 784 instr.addVariable("libs", "sharelib.keys", new Instrumentation.Variable<String>() { 785 @Override 786 public String getValue() { 787 Map<String, List<Path>> shareLib = getShareLib(); 788 if (shareLib != null && !shareLib.isEmpty()) { 789 Set<String> keySet = shareLib.keySet(); 790 return keySet.toString(); 791 } 792 return "(unavailable)"; 793 } 794 }); 795 instr.addVariable("libs", "launcherlib.system.libpath", new Instrumentation.Variable<String>() { 796 @Override 797 public String getValue() { 798 return getLauncherlibPath().toUri().toString(); 799 } 800 }); 801 instr.addVariable("libs", "sharelib.symlink.mapping", new Instrumentation.Variable<String>() { 802 @Override 803 public String getValue() { 804 Map<String, Map<Path, Path>> shareLibSymlinkMapping = getSymlinkMapping(); 805 if (shareLibSymlinkMapping != null && !shareLibSymlinkMapping.isEmpty() 806 && shareLibSymlinkMapping.values() != null && !shareLibSymlinkMapping.values().isEmpty()) { 807 StringBuffer bf = new StringBuffer(); 808 for (Entry<String, Map<Path, Path>> entry : shareLibSymlinkMapping.entrySet()) { 809 if (entry.getKey() != null && !entry.getValue().isEmpty()) { 810 for (Path path : entry.getValue().keySet()) { 811 bf.append(path).append("(").append(entry.getKey()).append(")").append("=>") 812 .append(shareLibSymlinkMapping.get(entry.getKey()) != null ? shareLibSymlinkMapping 813 .get(entry.getKey()).get(path) : "").append(","); 814 } 815 } 816 } 817 return bf.toString(); 818 } 819 return "(none)"; 820 } 821 }); 822 823 instr.addVariable("libs", "sharelib.cached.config.file", new Instrumentation.Variable<String>() { 824 @Override 825 public String getValue() { 826 Map<String, Map<Path, Configuration>> shareLibConfigMap = getShareLibConfigMap(); 827 if (shareLibConfigMap != null && !shareLibConfigMap.isEmpty()) { 828 StringBuffer bf = new StringBuffer(); 829 830 for (String path : shareLibConfigMap.keySet()) { 831 bf.append(path).append(";"); 832 } 833 return bf.toString(); 834 } 835 return "(none)"; 836 } 837 }); 838 839 } 840 841 /** 842 * Returns file system for shared libraries. 843 * <p/> 844 * If WorkflowAppService#getSystemLibPath doesn't have authority then a default one assumed 845 * 846 * @return file system for shared libraries 847 */ 848 public FileSystem getFileSystem() { 849 return fs; 850 } 851 852 /** 853 * Cache XML conf file 854 * 855 * @param hdfsPath the hdfs path 856 * @param shareLibKey the share lib key 857 * @throws IOException Signals that an I/O exception has occurred. 858 * @throws JDOMException 859 */ 860 private void cachePropertyFile(Path hdfsPath, String shareLibKey, 861 Map<String, Map<Path, Configuration>> shareLibConfigMap) throws IOException, JDOMException { 862 Map<Path, Configuration> confMap = shareLibConfigMap.get(shareLibKey); 863 if (confMap == null) { 864 confMap = new HashMap<Path, Configuration>(); 865 shareLibConfigMap.put(shareLibKey, confMap); 866 } 867 Configuration xmlConf = new XConfiguration(fs.open(hdfsPath)); 868 confMap.put(hdfsPath, xmlConf); 869 870 } 871 872 private void cacheActionKeySharelibConfList() { 873 ActionService actionService = Services.get().get(ActionService.class); 874 Set<String> actionTypes = actionService.getActionTypes(); 875 for (String key : actionTypes) { 876 ActionExecutor executor = actionService.getExecutor(key); 877 if (executor instanceof JavaActionExecutor) { 878 JavaActionExecutor jexecutor = (JavaActionExecutor) executor; 879 actionConfSet.addAll( 880 new HashSet<String>(Arrays.asList(jexecutor.getShareLibFilesForActionConf() == null ? new String[0] 881 : jexecutor.getShareLibFilesForActionConf()))); 882 } 883 } 884 } 885 886 public Configuration getShareLibConf(String inputKey, Path path) { 887 if (shareLibConfigMap.containsKey(inputKey)) { 888 return shareLibConfigMap.get(inputKey).get(path); 889 } 890 891 return null; 892 } 893 894 @VisibleForTesting 895 public Map<String, Map<Path, Configuration>> getShareLibConfigMap() { 896 return shareLibConfigMap; 897 } 898 899 private boolean isFilePartOfConfList(Path path) throws URISyntaxException { 900 String fragmentName = new URI(path.toString()).getFragment(); 901 String fileName = fragmentName == null ? path.getName() : fragmentName; 902 return actionConfSet.contains(fileName); 903 } 904}