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.client; 020 021import org.apache.oozie.BuildInfo; 022import org.apache.oozie.client.rest.JsonTags; 023import org.apache.oozie.client.rest.JsonToBean; 024import org.apache.oozie.client.rest.RestConstants; 025import org.apache.oozie.client.retry.ConnectionRetriableClient; 026import org.json.simple.JSONArray; 027import org.json.simple.JSONObject; 028import org.json.simple.JSONValue; 029import org.w3c.dom.Document; 030import org.w3c.dom.Element; 031 032import javax.xml.parsers.DocumentBuilderFactory; 033import javax.xml.transform.Transformer; 034import javax.xml.transform.TransformerFactory; 035import javax.xml.transform.dom.DOMSource; 036import javax.xml.transform.stream.StreamResult; 037import java.io.BufferedReader; 038import java.io.File; 039import java.io.FileInputStream; 040import java.io.IOException; 041import java.io.InputStream; 042import java.io.InputStreamReader; 043import java.io.OutputStream; 044import java.io.PrintStream; 045import java.io.Reader; 046import java.net.HttpURLConnection; 047import java.net.URL; 048import java.net.URLEncoder; 049import java.util.ArrayList; 050import java.util.Collections; 051import java.util.HashMap; 052import java.util.HashSet; 053import java.util.Iterator; 054import java.util.LinkedHashMap; 055import java.util.List; 056import java.util.Map; 057import java.util.Map.Entry; 058import java.util.Properties; 059import java.util.Set; 060import java.util.concurrent.Callable; 061 062/** 063 * Client API to submit and manage Oozie workflow jobs against an Oozie intance. 064 * <p/> 065 * This class is thread safe. 066 * <p/> 067 * Syntax for filter for the {@link #getJobsInfo(String)} {@link #getJobsInfo(String, int, int)} methods: 068 * <code>[NAME=VALUE][;NAME=VALUE]*</code>. 069 * <p/> 070 * Valid filter names are: 071 * <p/> 072 * <ul/> 073 * <li>name: the workflow application name from the workflow definition.</li> 074 * <li>user: the user that submitted the job.</li> 075 * <li>group: the group for the job.</li> 076 * <li>status: the status of the job.</li> 077 * </ul> 078 * <p/> 079 * The query will do an AND among all the filter names. The query will do an OR among all the filter values for the same 080 * name. Multiple values must be specified as different name value pairs. 081 */ 082public class OozieClient { 083 084 public static final long WS_PROTOCOL_VERSION_0 = 0; 085 086 public static final long WS_PROTOCOL_VERSION_1 = 1; 087 088 public static final long WS_PROTOCOL_VERSION = 2; // pointer to current version 089 090 public static final String USER_NAME = "user.name"; 091 092 @Deprecated 093 public static final String GROUP_NAME = "group.name"; 094 095 public static final String JOB_ACL = "oozie.job.acl"; 096 097 public static final String APP_PATH = "oozie.wf.application.path"; 098 099 public static final String COORDINATOR_APP_PATH = "oozie.coord.application.path"; 100 101 public static final String BUNDLE_APP_PATH = "oozie.bundle.application.path"; 102 103 public static final String BUNDLE_ID = "oozie.bundle.id"; 104 105 public static final String EXTERNAL_ID = "oozie.wf.external.id"; 106 107 public static final String WORKFLOW_NOTIFICATION_URL = "oozie.wf.workflow.notification.url"; 108 109 public static final String WORKFLOW_NOTIFICATION_PROXY = "oozie.wf.workflow.notification.proxy"; 110 111 public static final String ACTION_NOTIFICATION_URL = "oozie.wf.action.notification.url"; 112 113 public static final String COORD_ACTION_NOTIFICATION_URL = "oozie.coord.action.notification.url"; 114 115 public static final String COORD_ACTION_NOTIFICATION_PROXY = "oozie.coord.action.notification.proxy"; 116 117 public static final String RERUN_SKIP_NODES = "oozie.wf.rerun.skip.nodes"; 118 119 public static final String RERUN_FAIL_NODES = "oozie.wf.rerun.failnodes"; 120 121 public static final String LOG_TOKEN = "oozie.wf.log.token"; 122 123 public static final String ACTION_MAX_RETRIES = "oozie.wf.action.max.retries"; 124 125 public static final String ACTION_RETRY_INTERVAL = "oozie.wf.action.retry.interval"; 126 127 public static final String FILTER_USER = "user"; 128 129 public static final String FILTER_TEXT = "text"; 130 131 public static final String FILTER_GROUP = "group"; 132 133 public static final String FILTER_NAME = "name"; 134 135 public static final String FILTER_STATUS = "status"; 136 137 public static final String FILTER_NOMINAL_TIME = "nominaltime"; 138 139 public static final String FILTER_FREQUENCY = "frequency"; 140 141 public static final String FILTER_ID = "id"; 142 143 public static final String FILTER_UNIT = "unit"; 144 145 public static final String FILTER_JOBID = "jobid"; 146 147 public static final String FILTER_APPNAME = "appname"; 148 149 public static final String FILTER_SLA_APPNAME = "app_name"; 150 151 public static final String FILTER_SLA_ID = "id"; 152 153 public static final String FILTER_SLA_PARENT_ID = "parent_id"; 154 155 public static final String FILTER_SLA_NOMINAL_START = "nominal_start"; 156 157 public static final String FILTER_SLA_NOMINAL_END = "nominal_end"; 158 159 public static final String FILTER_CREATED_TIME_START = "startcreatedtime"; 160 161 public static final String FILTER_CREATED_TIME_END = "endcreatedtime"; 162 163 public static final String CHANGE_VALUE_ENDTIME = "endtime"; 164 165 public static final String CHANGE_VALUE_PAUSETIME = "pausetime"; 166 167 public static final String CHANGE_VALUE_CONCURRENCY = "concurrency"; 168 169 public static final String CHANGE_VALUE_STATUS = "status"; 170 171 public static final String LIBPATH = "oozie.libpath"; 172 173 public static final String USE_SYSTEM_LIBPATH = "oozie.use.system.libpath"; 174 175 public static final String OOZIE_SUSPEND_ON_NODES = "oozie.suspend.on.nodes"; 176 177 public static final String FILTER_SORT_BY = "sortby"; 178 179 public enum SORT_BY { 180 createdTime("createdTimestamp"), lastModifiedTime("lastModifiedTimestamp"); 181 private final String fullname; 182 183 SORT_BY(String fullname) { 184 this.fullname = fullname; 185 } 186 187 public String getFullname() { 188 return fullname; 189 } 190 } 191 192 public static enum SYSTEM_MODE { 193 NORMAL, NOWEBSERVICE, SAFEMODE 194 } 195 196 private static final Set<String> COMPLETED_WF_STATUSES = new HashSet<String>(); 197 private static final Set<String> COMPLETED_COORD_AND_BUNDLE_STATUSES = new HashSet<String>(); 198 private static final Set<String> COMPLETED_COORD_ACTION_STATUSES = new HashSet<String>(); 199 static { 200 COMPLETED_WF_STATUSES.add(WorkflowJob.Status.FAILED.toString()); 201 COMPLETED_WF_STATUSES.add(WorkflowJob.Status.KILLED.toString()); 202 COMPLETED_WF_STATUSES.add(WorkflowJob.Status.SUCCEEDED.toString()); 203 COMPLETED_COORD_AND_BUNDLE_STATUSES.add(Job.Status.FAILED.toString()); 204 COMPLETED_COORD_AND_BUNDLE_STATUSES.add(Job.Status.KILLED.toString()); 205 COMPLETED_COORD_AND_BUNDLE_STATUSES.add(Job.Status.SUCCEEDED.toString()); 206 COMPLETED_COORD_AND_BUNDLE_STATUSES.add(Job.Status.DONEWITHERROR.toString()); 207 COMPLETED_COORD_AND_BUNDLE_STATUSES.add(Job.Status.IGNORED.toString()); 208 COMPLETED_COORD_ACTION_STATUSES.add(CoordinatorAction.Status.FAILED.toString()); 209 COMPLETED_COORD_ACTION_STATUSES.add(CoordinatorAction.Status.IGNORED.toString()); 210 COMPLETED_COORD_ACTION_STATUSES.add(CoordinatorAction.Status.KILLED.toString()); 211 COMPLETED_COORD_ACTION_STATUSES.add(CoordinatorAction.Status.SKIPPED.toString()); 212 COMPLETED_COORD_ACTION_STATUSES.add(CoordinatorAction.Status.SUCCEEDED.toString()); 213 COMPLETED_COORD_ACTION_STATUSES.add(CoordinatorAction.Status.TIMEDOUT.toString()); 214 } 215 216 /** 217 * debugMode =0 means no debugging. > 0 means debugging on. 218 */ 219 public int debugMode = 0; 220 221 private int retryCount = 4; 222 223 224 private String baseUrl; 225 private String protocolUrl; 226 private boolean validatedVersion = false; 227 private JSONArray supportedVersions; 228 private final Map<String, String> headers = new HashMap<String, String>(); 229 230 private static final ThreadLocal<String> USER_NAME_TL = new ThreadLocal<String>(); 231 232 /** 233 * Allows to impersonate other users in the Oozie server. The current user 234 * must be configured as a proxyuser in Oozie. 235 * <p/> 236 * IMPORTANT: impersonation happens only with Oozie client requests done within 237 * doAs() calls. 238 * 239 * @param userName user to impersonate. 240 * @param callable callable with {@link OozieClient} calls impersonating the specified user. 241 * @return any response returned by the {@link Callable#call()} method. 242 * @throws Exception thrown by the {@link Callable#call()} method. 243 */ 244 public static <T> T doAs(String userName, Callable<T> callable) throws Exception { 245 notEmpty(userName, "userName"); 246 notNull(callable, "callable"); 247 try { 248 USER_NAME_TL.set(userName); 249 return callable.call(); 250 } 251 finally { 252 USER_NAME_TL.remove(); 253 } 254 } 255 256 protected OozieClient() { 257 } 258 259 /** 260 * Create a Workflow client instance. 261 * 262 * @param oozieUrl URL of the Oozie instance it will interact with. 263 */ 264 public OozieClient(String oozieUrl) { 265 this.baseUrl = notEmpty(oozieUrl, "oozieUrl"); 266 if (!this.baseUrl.endsWith("/")) { 267 this.baseUrl += "/"; 268 } 269 } 270 271 /** 272 * Return the Oozie URL of the workflow client instance. 273 * <p/> 274 * This URL is the base URL fo the Oozie system, with not protocol versioning. 275 * 276 * @return the Oozie URL of the workflow client instance. 277 */ 278 public String getOozieUrl() { 279 return baseUrl; 280 } 281 282 /** 283 * Return the Oozie URL used by the client and server for WS communications. 284 * <p/> 285 * This URL is the original URL plus the versioning element path. 286 * 287 * @return the Oozie URL used by the client and server for communication. 288 * @throws OozieClientException thrown in the client and the server are not protocol compatible. 289 */ 290 public String getProtocolUrl() throws OozieClientException { 291 validateWSVersion(); 292 return protocolUrl; 293 } 294 295 /** 296 * @return current debug Mode 297 */ 298 public int getDebugMode() { 299 return debugMode; 300 } 301 302 /** 303 * Set debug mode. 304 * 305 * @param debugMode : 0 means no debugging. > 0 means debugging 306 */ 307 public void setDebugMode(int debugMode) { 308 this.debugMode = debugMode; 309 } 310 311 public int getRetryCount() { 312 return retryCount; 313 } 314 315 316 public void setRetryCount(int retryCount) { 317 this.retryCount = retryCount; 318 } 319 320 private String getBaseURLForVersion(long protocolVersion) throws OozieClientException { 321 try { 322 if (supportedVersions == null) { 323 supportedVersions = getSupportedProtocolVersions(); 324 } 325 if (supportedVersions == null) { 326 throw new OozieClientException("HTTP error", "no response message"); 327 } 328 if (supportedVersions.contains(protocolVersion)) { 329 return baseUrl + "v" + protocolVersion + "/"; 330 } 331 else { 332 throw new OozieClientException(OozieClientException.UNSUPPORTED_VERSION, "Protocol version " 333 + protocolVersion + " is not supported"); 334 } 335 } 336 catch (IOException e) { 337 throw new OozieClientException(OozieClientException.IO_ERROR, e); 338 } 339 } 340 341 /** 342 * Validate that the Oozie client and server instances are protocol compatible. 343 * 344 * @throws OozieClientException thrown in the client and the server are not protocol compatible. 345 */ 346 public synchronized void validateWSVersion() throws OozieClientException { 347 if (!validatedVersion) { 348 try { 349 supportedVersions = getSupportedProtocolVersions(); 350 if (supportedVersions == null) { 351 throw new OozieClientException("HTTP error", "no response message"); 352 } 353 if (!supportedVersions.contains(WS_PROTOCOL_VERSION) 354 && !supportedVersions.contains(WS_PROTOCOL_VERSION_1) 355 && !supportedVersions.contains(WS_PROTOCOL_VERSION_0)) { 356 StringBuilder msg = new StringBuilder(); 357 msg.append("Supported version [").append(WS_PROTOCOL_VERSION) 358 .append("] or less, Unsupported versions["); 359 String separator = ""; 360 for (Object version : supportedVersions) { 361 msg.append(separator).append(version); 362 } 363 msg.append("]"); 364 throw new OozieClientException(OozieClientException.UNSUPPORTED_VERSION, msg.toString()); 365 } 366 if (supportedVersions.contains(WS_PROTOCOL_VERSION)) { 367 protocolUrl = baseUrl + "v" + WS_PROTOCOL_VERSION + "/"; 368 } 369 else if (supportedVersions.contains(WS_PROTOCOL_VERSION_1)) { 370 protocolUrl = baseUrl + "v" + WS_PROTOCOL_VERSION_1 + "/"; 371 } 372 else { 373 if (supportedVersions.contains(WS_PROTOCOL_VERSION_0)) { 374 protocolUrl = baseUrl + "v" + WS_PROTOCOL_VERSION_0 + "/"; 375 } 376 } 377 } 378 catch (IOException ex) { 379 throw new OozieClientException(OozieClientException.IO_ERROR, ex); 380 } 381 validatedVersion = true; 382 } 383 } 384 385 private JSONArray getSupportedProtocolVersions() throws IOException, OozieClientException { 386 JSONArray versions = null; 387 final URL url = new URL(baseUrl + RestConstants.VERSIONS); 388 389 HttpURLConnection conn = createRetryableConnection(url, "GET"); 390 391 if (conn.getResponseCode() == HttpURLConnection.HTTP_OK) { 392 versions = (JSONArray) JSONValue.parse(new InputStreamReader(conn.getInputStream())); 393 } 394 else { 395 handleError(conn); 396 } 397 return versions; 398 } 399 400 /** 401 * Create an empty configuration with just the {@link #USER_NAME} set to the JVM user name. 402 * 403 * @return an empty configuration. 404 */ 405 public Properties createConfiguration() { 406 Properties conf = new Properties(); 407 String userName = USER_NAME_TL.get(); 408 if (userName == null) { 409 userName = System.getProperty("user.name"); 410 } 411 conf.setProperty(USER_NAME, userName); 412 return conf; 413 } 414 415 /** 416 * Set a HTTP header to be used in the WS requests by the workflow instance. 417 * 418 * @param name header name. 419 * @param value header value. 420 */ 421 public void setHeader(String name, String value) { 422 headers.put(notEmpty(name, "name"), notNull(value, "value")); 423 } 424 425 /** 426 * Get the value of a set HTTP header from the workflow instance. 427 * 428 * @param name header name. 429 * @return header value, <code>null</code> if not set. 430 */ 431 public String getHeader(String name) { 432 return headers.get(notEmpty(name, "name")); 433 } 434 435 /** 436 * Get the set HTTP header 437 * 438 * @return map of header key and value 439 */ 440 public Map<String, String> getHeaders() { 441 return headers; 442 } 443 444 /** 445 * Remove a HTTP header from the workflow client instance. 446 * 447 * @param name header name. 448 */ 449 public void removeHeader(String name) { 450 headers.remove(notEmpty(name, "name")); 451 } 452 453 /** 454 * Return an iterator with all the header names set in the workflow instance. 455 * 456 * @return header names. 457 */ 458 public Iterator<String> getHeaderNames() { 459 return Collections.unmodifiableMap(headers).keySet().iterator(); 460 } 461 462 private URL createURL(Long protocolVersion, String collection, String resource, Map<String, String> parameters) 463 throws IOException, OozieClientException { 464 validateWSVersion(); 465 StringBuilder sb = new StringBuilder(); 466 if (protocolVersion == null) { 467 sb.append(protocolUrl); 468 } 469 else { 470 sb.append(getBaseURLForVersion(protocolVersion)); 471 } 472 sb.append(collection); 473 if (resource != null && resource.length() > 0) { 474 sb.append("/").append(resource); 475 } 476 if (parameters.size() > 0) { 477 String separator = "?"; 478 for (Map.Entry<String, String> param : parameters.entrySet()) { 479 if (param.getValue() != null) { 480 sb.append(separator).append(URLEncoder.encode(param.getKey(), "UTF-8")).append("=").append( 481 URLEncoder.encode(param.getValue(), "UTF-8")); 482 separator = "&"; 483 } 484 } 485 } 486 return new URL(sb.toString()); 487 } 488 489 private boolean validateCommand(String url) { 490 { 491 if (protocolUrl.contains(baseUrl + "v0")) { 492 if (url.contains("dryrun") || url.contains("jobtype=c") || url.contains("systemmode")) { 493 return false; 494 } 495 } 496 } 497 return true; 498 } 499 /** 500 * Create retryable http connection to oozie server. 501 * 502 * @param url 503 * @param method 504 * @return connection 505 * @throws IOException 506 * @throws OozieClientException 507 */ 508 protected HttpURLConnection createRetryableConnection(final URL url, final String method) throws IOException{ 509 return (HttpURLConnection) new ConnectionRetriableClient(getRetryCount()) { 510 @Override 511 public Object doExecute(URL url, String method) throws IOException, OozieClientException { 512 HttpURLConnection conn = createConnection(url, method); 513 return conn; 514 } 515 }.execute(url, method); 516 } 517 518 /** 519 * Create http connection to oozie server. 520 * 521 * @param url 522 * @param method 523 * @return connection 524 * @throws IOException 525 * @throws OozieClientException 526 */ 527 protected HttpURLConnection createConnection(URL url, String method) throws IOException, OozieClientException { 528 HttpURLConnection conn = (HttpURLConnection) url.openConnection(); 529 conn.setRequestMethod(method); 530 if (method.equals("POST") || method.equals("PUT")) { 531 conn.setDoOutput(true); 532 } 533 for (Map.Entry<String, String> header : headers.entrySet()) { 534 conn.setRequestProperty(header.getKey(), header.getValue()); 535 } 536 return conn; 537 } 538 539 protected abstract class ClientCallable<T> implements Callable<T> { 540 private final String method; 541 private final String collection; 542 private final String resource; 543 private final Map<String, String> params; 544 private final Long protocolVersion; 545 546 public ClientCallable(String method, String collection, String resource, Map<String, String> params) { 547 this(method, null, collection, resource, params); 548 } 549 550 public ClientCallable(String method, Long protocolVersion, String collection, String resource, Map<String, String> params) { 551 this.method = method; 552 this.protocolVersion = protocolVersion; 553 this.collection = collection; 554 this.resource = resource; 555 this.params = params; 556 } 557 558 public T call() throws OozieClientException { 559 try { 560 URL url = createURL(protocolVersion, collection, resource, params); 561 if (validateCommand(url.toString())) { 562 if (getDebugMode() > 0) { 563 System.out.println(method + " " + url); 564 } 565 return call(createRetryableConnection(url, method)); 566 } 567 else { 568 System.out.println("Option not supported in target server. Supported only on Oozie-2.0 or greater." 569 + " Use 'oozie help' for details"); 570 throw new OozieClientException(OozieClientException.UNSUPPORTED_VERSION, new Exception()); 571 } 572 } 573 catch (IOException ex) { 574 throw new OozieClientException(OozieClientException.IO_ERROR, ex); 575 } 576 } 577 578 protected abstract T call(HttpURLConnection conn) throws IOException, OozieClientException; 579 } 580 581 protected abstract class MapClientCallable extends ClientCallable<Map<String, String>> { 582 583 MapClientCallable(String method, String collection, String resource, Map<String, String> params) { 584 super(method, collection, resource, params); 585 } 586 587 @Override 588 protected Map<String, String> call(HttpURLConnection conn) throws IOException, OozieClientException { 589 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 590 Reader reader = new InputStreamReader(conn.getInputStream()); 591 JSONObject json = (JSONObject) JSONValue.parse(reader); 592 Map<String, String> map = new HashMap<String, String>(); 593 for (Object key : json.keySet()) { 594 map.put((String)key, (String)json.get(key)); 595 } 596 return map; 597 } 598 else { 599 handleError(conn); 600 } 601 return null; 602 } 603 } 604 605 static void handleError(HttpURLConnection conn) throws IOException, OozieClientException { 606 int status = conn.getResponseCode(); 607 String error = conn.getHeaderField(RestConstants.OOZIE_ERROR_CODE); 608 String message = conn.getHeaderField(RestConstants.OOZIE_ERROR_MESSAGE); 609 610 if (error == null) { 611 error = "HTTP error code: " + status; 612 } 613 614 if (message == null) { 615 message = conn.getResponseMessage(); 616 } 617 throw new OozieClientException(error, message); 618 } 619 620 static Map<String, String> prepareParams(String... params) { 621 Map<String, String> map = new LinkedHashMap<String, String>(); 622 for (int i = 0; i < params.length; i = i + 2) { 623 map.put(params[i], params[i + 1]); 624 } 625 String doAsUserName = USER_NAME_TL.get(); 626 if (doAsUserName != null) { 627 map.put(RestConstants.DO_AS_PARAM, doAsUserName); 628 } 629 return map; 630 } 631 632 public void writeToXml(Properties props, OutputStream out) throws IOException { 633 try { 634 Document doc = DocumentBuilderFactory.newInstance().newDocumentBuilder().newDocument(); 635 Element conf = doc.createElement("configuration"); 636 doc.appendChild(conf); 637 conf.appendChild(doc.createTextNode("\n")); 638 for (String name : props.stringPropertyNames()) { // Properties whose key or value is not of type String are omitted. 639 String value = props.getProperty(name); 640 Element propNode = doc.createElement("property"); 641 conf.appendChild(propNode); 642 643 Element nameNode = doc.createElement("name"); 644 nameNode.appendChild(doc.createTextNode(name.trim())); 645 propNode.appendChild(nameNode); 646 647 Element valueNode = doc.createElement("value"); 648 valueNode.appendChild(doc.createTextNode(value.trim())); 649 propNode.appendChild(valueNode); 650 651 conf.appendChild(doc.createTextNode("\n")); 652 } 653 654 DOMSource source = new DOMSource(doc); 655 StreamResult result = new StreamResult(out); 656 TransformerFactory transFactory = TransformerFactory.newInstance(); 657 transFactory.setFeature("http://javax.xml.XMLConstants/feature/secure-processing", true); 658 Transformer transformer = transFactory.newTransformer(); 659 transformer.transform(source, result); 660 if (getDebugMode() > 0) { 661 result = new StreamResult(System.out); 662 transformer.transform(source, result); 663 System.out.println(); 664 } 665 } 666 catch (Exception e) { 667 throw new IOException(e); 668 } 669 } 670 671 private class JobSubmit extends ClientCallable<String> { 672 private final Properties conf; 673 674 JobSubmit(Properties conf, boolean start) { 675 super("POST", RestConstants.JOBS, "", (start) ? prepareParams(RestConstants.ACTION_PARAM, 676 RestConstants.JOB_ACTION_START) : prepareParams()); 677 this.conf = notNull(conf, "conf"); 678 } 679 680 JobSubmit(String jobId, Properties conf) { 681 super("PUT", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.ACTION_PARAM, 682 RestConstants.JOB_ACTION_RERUN)); 683 this.conf = notNull(conf, "conf"); 684 } 685 686 public JobSubmit(Properties conf, String jobActionDryrun) { 687 super("POST", RestConstants.JOBS, "", prepareParams(RestConstants.ACTION_PARAM, 688 RestConstants.JOB_ACTION_DRYRUN)); 689 this.conf = notNull(conf, "conf"); 690 } 691 692 @Override 693 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 694 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 695 writeToXml(conf, conn.getOutputStream()); 696 if (conn.getResponseCode() == HttpURLConnection.HTTP_CREATED) { 697 JSONObject json = (JSONObject) JSONValue.parse(new InputStreamReader(conn.getInputStream())); 698 return (String) json.get(JsonTags.JOB_ID); 699 } 700 if (conn.getResponseCode() != HttpURLConnection.HTTP_OK) { 701 handleError(conn); 702 } 703 return null; 704 } 705 } 706 707 /** 708 * Submit a workflow job. 709 * 710 * @param conf job configuration. 711 * @return the job Id. 712 * @throws OozieClientException thrown if the job could not be submitted. 713 */ 714 public String submit(Properties conf) throws OozieClientException { 715 return (new JobSubmit(conf, false)).call(); 716 } 717 718 private class JobAction extends ClientCallable<Void> { 719 720 JobAction(String jobId, String action) { 721 super("PUT", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.ACTION_PARAM, action)); 722 } 723 724 JobAction(String jobId, String action, String params) { 725 super("PUT", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.ACTION_PARAM, action, 726 RestConstants.JOB_CHANGE_VALUE, params)); 727 } 728 729 @Override 730 protected Void call(HttpURLConnection conn) throws IOException, OozieClientException { 731 if (!(conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 732 handleError(conn); 733 } 734 return null; 735 } 736 } 737 738 private class JobsAction extends ClientCallable<JSONObject> { 739 740 JobsAction(String action, String filter, String jobType, int start, int len) { 741 super("PUT", RestConstants.JOBS, "", 742 prepareParams(RestConstants.ACTION_PARAM, action, 743 RestConstants.JOB_FILTER_PARAM, filter, RestConstants.JOBTYPE_PARAM, jobType, 744 RestConstants.OFFSET_PARAM, Integer.toString(start), 745 RestConstants.LEN_PARAM, Integer.toString(len))); 746 } 747 748 @Override 749 protected JSONObject call(HttpURLConnection conn) throws IOException, OozieClientException { 750 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 751 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 752 Reader reader = new InputStreamReader(conn.getInputStream()); 753 JSONObject json = (JSONObject) JSONValue.parse(reader); 754 return json; 755 } 756 else { 757 handleError(conn); 758 } 759 return null; 760 } 761 } 762 /** 763 * Update coord definition. 764 * 765 * @param jobId the job id 766 * @param conf the conf 767 * @param dryrun the dryrun 768 * @param showDiff the show diff 769 * @return the string 770 * @throws OozieClientException the oozie client exception 771 */ 772 public String updateCoord(String jobId, Properties conf, String dryrun, String showDiff) 773 throws OozieClientException { 774 return (new UpdateCoord(jobId, conf, dryrun, showDiff)).call(); 775 } 776 777 /** 778 * Update coord definition without properties. 779 * 780 * @param jobId the job id 781 * @param dryrun the dryrun 782 * @param showDiff the show diff 783 * @return the string 784 * @throws OozieClientException the oozie client exception 785 */ 786 public String updateCoord(String jobId, String dryrun, String showDiff) throws OozieClientException { 787 return (new UpdateCoord(jobId, dryrun, showDiff)).call(); 788 } 789 790 /** 791 * The Class UpdateCoord. 792 */ 793 private class UpdateCoord extends ClientCallable<String> { 794 private final Properties conf; 795 796 public UpdateCoord(String jobId, Properties conf, String jobActionDryrun, String showDiff) { 797 super("PUT", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.ACTION_PARAM, 798 RestConstants.JOB_COORD_UPDATE, RestConstants.JOB_ACTION_DRYRUN, jobActionDryrun, 799 RestConstants.JOB_ACTION_SHOWDIFF, showDiff)); 800 this.conf = conf; 801 } 802 803 public UpdateCoord(String jobId, String jobActionDryrun, String showDiff) { 804 this(jobId, new Properties(), jobActionDryrun, showDiff); 805 } 806 807 @Override 808 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 809 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 810 writeToXml(conf, conn.getOutputStream()); 811 812 if (conn.getResponseCode() == HttpURLConnection.HTTP_OK) { 813 JSONObject json = (JSONObject) JSONValue.parse(new InputStreamReader(conn.getInputStream())); 814 JSONObject update = (JSONObject) json.get(JsonTags.COORD_UPDATE); 815 if (update != null) { 816 return (String) update.get(JsonTags.COORD_UPDATE_DIFF); 817 } 818 else { 819 return ""; 820 } 821 } 822 if (conn.getResponseCode() != HttpURLConnection.HTTP_OK) { 823 handleError(conn); 824 } 825 return null; 826 } 827 } 828 829 /** 830 * dryrun for a given job 831 * 832 * @param conf Job configuration. 833 */ 834 public String dryrun(Properties conf) throws OozieClientException { 835 return new JobSubmit(conf, RestConstants.JOB_ACTION_DRYRUN).call(); 836 } 837 838 /** 839 * Start a workflow job. 840 * 841 * @param jobId job Id. 842 * @throws OozieClientException thrown if the job could not be started. 843 */ 844 public void start(String jobId) throws OozieClientException { 845 new JobAction(jobId, RestConstants.JOB_ACTION_START).call(); 846 } 847 848 /** 849 * Submit and start a workflow job. 850 * 851 * @param conf job configuration. 852 * @return the job Id. 853 * @throws OozieClientException thrown if the job could not be submitted. 854 */ 855 public String run(Properties conf) throws OozieClientException { 856 return (new JobSubmit(conf, true)).call(); 857 } 858 859 /** 860 * Rerun a workflow job. 861 * 862 * @param jobId job Id to rerun. 863 * @param conf configuration information for the rerun. 864 * @throws OozieClientException thrown if the job could not be started. 865 */ 866 public void reRun(String jobId, Properties conf) throws OozieClientException { 867 new JobSubmit(jobId, conf).call(); 868 } 869 870 /** 871 * Suspend a workflow job. 872 * 873 * @param jobId job Id. 874 * @throws OozieClientException thrown if the job could not be suspended. 875 */ 876 public void suspend(String jobId) throws OozieClientException { 877 new JobAction(jobId, RestConstants.JOB_ACTION_SUSPEND).call(); 878 } 879 880 /** 881 * Resume a workflow job. 882 * 883 * @param jobId job Id. 884 * @throws OozieClientException thrown if the job could not be resume. 885 */ 886 public void resume(String jobId) throws OozieClientException { 887 new JobAction(jobId, RestConstants.JOB_ACTION_RESUME).call(); 888 } 889 890 /** 891 * Kill a workflow/coord/bundle job. 892 * 893 * @param jobId job Id. 894 * @throws OozieClientException thrown if the job could not be killed. 895 */ 896 public void kill(String jobId) throws OozieClientException { 897 new JobAction(jobId, RestConstants.JOB_ACTION_KILL).call(); 898 } 899 900 /** 901 * Kill coordinator actions 902 * @param jobId coordinator Job Id 903 * @param rangeType type 'date' if -date is used, 'action-num' if -action is used 904 * @param scope kill scope for date or action nums 905 * @return list of coordinator actions that underwent kill 906 * @throws OozieClientException thrown if some actions could not be killed. 907 */ 908 public List<CoordinatorAction> kill(String jobId, String rangeType, String scope) throws OozieClientException { 909 return new CoordActionsKill(jobId, rangeType, scope).call(); 910 } 911 912 public JSONObject bulkModifyJobs(String actionType, String filter, String jobType, int start, int len) 913 throws OozieClientException { 914 return new JobsAction(actionType, filter, jobType, start, len).call(); 915 } 916 917 public JSONObject killJobs(String filter, String jobType, int start, int len) 918 throws OozieClientException { 919 return bulkModifyJobs("kill", filter, jobType, start, len); 920 } 921 922 public JSONObject suspendJobs(String filter, String jobType, int start, int len) 923 throws OozieClientException { 924 return bulkModifyJobs("suspend", filter, jobType, start, len); 925 } 926 927 public JSONObject resumeJobs(String filter, String jobType, int start, int len) 928 throws OozieClientException { 929 return bulkModifyJobs("resume", filter, jobType, start, len); 930 } 931 /** 932 * Change a coordinator job. 933 * 934 * @param jobId job Id. 935 * @param changeValue change value. 936 * @throws OozieClientException thrown if the job could not be changed. 937 */ 938 public void change(String jobId, String changeValue) throws OozieClientException { 939 new JobAction(jobId, RestConstants.JOB_ACTION_CHANGE, changeValue).call(); 940 } 941 942 /** 943 * Ignore a coordinator job. 944 * 945 * @param jobId coord job Id. 946 * @param scope list of coord actions to be ignored 947 * @throws OozieClientException thrown if the job could not be changed. 948 */ 949 public List<CoordinatorAction> ignore(String jobId, String scope) throws OozieClientException { 950 return new CoordIgnore(jobId, RestConstants.JOB_COORD_SCOPE_ACTION, scope).call(); 951 } 952 953 private class JobInfo extends ClientCallable<WorkflowJob> { 954 955 JobInfo(String jobId, int start, int len) { 956 super("GET", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.JOB_SHOW_PARAM, 957 RestConstants.JOB_SHOW_INFO, RestConstants.OFFSET_PARAM, Integer.toString(start), 958 RestConstants.LEN_PARAM, Integer.toString(len))); 959 } 960 961 @Override 962 protected WorkflowJob call(HttpURLConnection conn) throws IOException, OozieClientException { 963 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 964 Reader reader = new InputStreamReader(conn.getInputStream()); 965 JSONObject json = (JSONObject) JSONValue.parse(reader); 966 return JsonToBean.createWorkflowJob(json); 967 } 968 else { 969 handleError(conn); 970 } 971 return null; 972 } 973 } 974 975 private class JMSInfo extends ClientCallable<JMSConnectionInfo> { 976 977 JMSInfo() { 978 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_JMS_INFO, prepareParams()); 979 } 980 981 protected JMSConnectionInfo call(HttpURLConnection conn) throws IOException, OozieClientException { 982 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 983 Reader reader = new InputStreamReader(conn.getInputStream()); 984 JSONObject json = (JSONObject) JSONValue.parse(reader); 985 return JsonToBean.createJMSConnectionInfo(json); 986 } 987 else { 988 handleError(conn); 989 } 990 return null; 991 } 992 } 993 994 private class WorkflowActionInfo extends ClientCallable<WorkflowAction> { 995 WorkflowActionInfo(String actionId) { 996 super("GET", RestConstants.JOB, notEmpty(actionId, "id"), prepareParams(RestConstants.JOB_SHOW_PARAM, 997 RestConstants.JOB_SHOW_INFO)); 998 } 999 1000 @Override 1001 protected WorkflowAction call(HttpURLConnection conn) throws IOException, OozieClientException { 1002 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1003 Reader reader = new InputStreamReader(conn.getInputStream()); 1004 JSONObject json = (JSONObject) JSONValue.parse(reader); 1005 return JsonToBean.createWorkflowAction(json); 1006 } 1007 else { 1008 handleError(conn); 1009 } 1010 return null; 1011 } 1012 } 1013 1014 /** 1015 * Get the info of a workflow job. 1016 * 1017 * @param jobId job Id. 1018 * @return the job info. 1019 * @throws OozieClientException thrown if the job info could not be retrieved. 1020 */ 1021 public WorkflowJob getJobInfo(String jobId) throws OozieClientException { 1022 return getJobInfo(jobId, 0, 0); 1023 } 1024 1025 /** 1026 * Get the JMS Connection info 1027 * @return JMSConnectionInfo object 1028 * @throws OozieClientException 1029 */ 1030 public JMSConnectionInfo getJMSConnectionInfo() throws OozieClientException { 1031 return new JMSInfo().call(); 1032 } 1033 1034 /** 1035 * Get the info of a workflow job and subset actions. 1036 * 1037 * @param jobId job Id. 1038 * @param start starting index in the list of actions belonging to the job 1039 * @param len number of actions to be returned 1040 * @return the job info. 1041 * @throws OozieClientException thrown if the job info could not be retrieved. 1042 */ 1043 public WorkflowJob getJobInfo(String jobId, int start, int len) throws OozieClientException { 1044 return new JobInfo(jobId, start, len).call(); 1045 } 1046 1047 /** 1048 * Get the info of a workflow action. 1049 * 1050 * @param actionId Id. 1051 * @return the workflow action info. 1052 * @throws OozieClientException thrown if the job info could not be retrieved. 1053 */ 1054 public WorkflowAction getWorkflowActionInfo(String actionId) throws OozieClientException { 1055 return new WorkflowActionInfo(actionId).call(); 1056 } 1057 1058 /** 1059 * Get the log of a workflow job. 1060 * 1061 * @param jobId job Id. 1062 * @return the job log. 1063 * @throws OozieClientException thrown if the job info could not be retrieved. 1064 */ 1065 public String getJobLog(String jobId) throws OozieClientException { 1066 return new JobLog(jobId).call(); 1067 } 1068 1069 /** 1070 * Get the log of a job. 1071 * 1072 * @param jobId job Id. 1073 * @param logRetrievalType Based on which filter criteria the log is retrieved 1074 * @param logRetrievalScope Value for the retrieval type 1075 * @param logFilter log filter 1076 * @param ps Printstream of command line interface 1077 * @throws OozieClientException thrown if the job info could not be retrieved. 1078 */ 1079 public void getJobLog(String jobId, String logRetrievalType, String logRetrievalScope, String logFilter, 1080 PrintStream ps) throws OozieClientException { 1081 new JobLog(jobId, logRetrievalType, logRetrievalScope, logFilter, ps).call(); 1082 } 1083 1084 /** 1085 * Get the log of a job. 1086 * 1087 * @param jobId job Id. 1088 * @param logRetrievalType Based on which filter criteria the log is retrieved 1089 * @param logRetrievalScope Value for the retrieval type 1090 * @param ps Printstream of command line interface 1091 * @throws OozieClientException thrown if the job info could not be retrieved. 1092 */ 1093 public void getJobLog(String jobId, String logRetrievalType, String logRetrievalScope, PrintStream ps) 1094 throws OozieClientException { 1095 getJobLog(jobId, logRetrievalType, logRetrievalScope, null, ps); 1096 } 1097 1098 private class JobLog extends JobMetadata { 1099 JobLog(String jobId) { 1100 super(jobId, RestConstants.JOB_SHOW_LOG); 1101 } 1102 JobLog(String jobId, String logRetrievalType, String logRetrievalScope, String logFilter, PrintStream ps) { 1103 super(jobId, logRetrievalType, logRetrievalScope, RestConstants.JOB_SHOW_LOG, logFilter, ps); 1104 } 1105 } 1106 1107 /** 1108 * Gets the JMS topic name for a particular job 1109 * @param jobId given jobId 1110 * @return the JMS topic name 1111 * @throws OozieClientException 1112 */ 1113 public String getJMSTopicName(String jobId) throws OozieClientException { 1114 return new JMSTopic(jobId).call(); 1115 } 1116 1117 private class JMSTopic extends ClientCallable<String> { 1118 1119 JMSTopic(String jobId) { 1120 super("GET", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.JOB_SHOW_PARAM, 1121 RestConstants.JOB_SHOW_JMS_TOPIC)); 1122 } 1123 1124 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 1125 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1126 Reader reader = new InputStreamReader(conn.getInputStream()); 1127 JSONObject json = (JSONObject) JSONValue.parse(reader); 1128 return (String) json.get(JsonTags.JMS_TOPIC_NAME); 1129 } 1130 else { 1131 handleError(conn); 1132 } 1133 return null; 1134 } 1135 } 1136 1137 /** 1138 * Get the definition of a workflow job. 1139 * 1140 * @param jobId job Id. 1141 * @return the job log. 1142 * @throws OozieClientException thrown if the job info could not be retrieved. 1143 */ 1144 public String getJobDefinition(String jobId) throws OozieClientException { 1145 return new JobDefinition(jobId).call(); 1146 } 1147 1148 private class JobDefinition extends JobMetadata { 1149 1150 JobDefinition(String jobId) { 1151 super(jobId, RestConstants.JOB_SHOW_DEFINITION); 1152 } 1153 } 1154 1155 private class JobMetadata extends ClientCallable<String> { 1156 PrintStream printStream; 1157 1158 JobMetadata(String jobId, String metaType) { 1159 super("GET", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.JOB_SHOW_PARAM, 1160 metaType)); 1161 } 1162 1163 JobMetadata(String jobId, String logRetrievalType, String logRetrievalScope, String metaType, String logFilter, 1164 PrintStream ps) { 1165 super("GET", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.JOB_SHOW_PARAM, 1166 metaType, RestConstants.JOB_LOG_TYPE_PARAM, logRetrievalType, RestConstants.JOB_LOG_SCOPE_PARAM, 1167 logRetrievalScope, RestConstants.LOG_FILTER_OPTION, logFilter)); 1168 printStream = ps; 1169 } 1170 1171 @Override 1172 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 1173 String returnVal = null; 1174 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1175 InputStream is = conn.getInputStream(); 1176 InputStreamReader isr = new InputStreamReader(is); 1177 try { 1178 if (printStream != null) { 1179 sendToOutputStream(isr, -1); 1180 } 1181 else { 1182 returnVal = getReaderAsString(isr, -1); 1183 } 1184 } 1185 finally { 1186 isr.close(); 1187 } 1188 } 1189 else { 1190 handleError(conn); 1191 } 1192 return returnVal; 1193 } 1194 1195 /** 1196 * Output the log to command line interface 1197 * 1198 * @param reader reader to read into a string. 1199 * @param maxLen max content length allowed, if -1 there is no limit. 1200 * @throws IOException 1201 */ 1202 private void sendToOutputStream(Reader reader, int maxLen) throws IOException { 1203 if (reader == null) { 1204 throw new IllegalArgumentException("reader cannot be null"); 1205 } 1206 StringBuilder sb = new StringBuilder(); 1207 char[] buffer = new char[2048]; 1208 int read; 1209 int count = 0; 1210 int noOfCharstoFlush = 1024; 1211 while ((read = reader.read(buffer)) > -1) { 1212 count += read; 1213 if ((maxLen > -1) && (count > maxLen)) { 1214 break; 1215 } 1216 sb.append(buffer, 0, read); 1217 if (sb.length() > noOfCharstoFlush) { 1218 printStream.print(sb.toString()); 1219 sb = new StringBuilder(""); 1220 } 1221 } 1222 printStream.print(sb.toString()); 1223 } 1224 1225 /** 1226 * Return a reader as string. 1227 * <p/> 1228 * 1229 * @param reader reader to read into a string. 1230 * @param maxLen max content length allowed, if -1 there is no limit. 1231 * @return the reader content. 1232 * @throws IOException thrown if the resource could not be read. 1233 */ 1234 private String getReaderAsString(Reader reader, int maxLen) throws IOException { 1235 if (reader == null) { 1236 throw new IllegalArgumentException("reader cannot be null"); 1237 } 1238 StringBuffer sb = new StringBuffer(); 1239 char[] buffer = new char[2048]; 1240 int read; 1241 int count = 0; 1242 while ((read = reader.read(buffer)) > -1) { 1243 count += read; 1244 1245 // read up to maxLen chars; 1246 if ((maxLen > -1) && (count > maxLen)) { 1247 break; 1248 } 1249 sb.append(buffer, 0, read); 1250 } 1251 reader.close(); 1252 return sb.toString(); 1253 } 1254 } 1255 1256 private class CoordJobInfo extends ClientCallable<CoordinatorJob> { 1257 1258 CoordJobInfo(String jobId, String filter, int start, int len, String order) { 1259 super("GET", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.JOB_SHOW_PARAM, 1260 RestConstants.JOB_SHOW_INFO, RestConstants.JOB_FILTER_PARAM, filter, RestConstants.OFFSET_PARAM, 1261 Integer.toString(start), RestConstants.LEN_PARAM, Integer.toString(len), RestConstants.ORDER_PARAM, 1262 order)); 1263 } 1264 1265 @Override 1266 protected CoordinatorJob call(HttpURLConnection conn) throws IOException, OozieClientException { 1267 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1268 Reader reader = new InputStreamReader(conn.getInputStream()); 1269 JSONObject json = (JSONObject) JSONValue.parse(reader); 1270 return JsonToBean.createCoordinatorJob(json); 1271 } 1272 else { 1273 handleError(conn); 1274 } 1275 return null; 1276 } 1277 } 1278 1279 private class WfsForCoordAction extends ClientCallable<List<WorkflowJob>> { 1280 1281 WfsForCoordAction(String coordActionId) { 1282 super("GET", RestConstants.JOB, notEmpty(coordActionId, "coordActionId"), prepareParams( 1283 RestConstants.JOB_SHOW_PARAM, RestConstants.ALL_WORKFLOWS_FOR_COORD_ACTION)); 1284 } 1285 1286 @Override 1287 protected List<WorkflowJob> call(HttpURLConnection conn) throws IOException, OozieClientException { 1288 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1289 Reader reader = new InputStreamReader(conn.getInputStream()); 1290 JSONObject json = (JSONObject) JSONValue.parse(reader); 1291 JSONArray workflows = (JSONArray) json.get(JsonTags.WORKFLOWS_JOBS); 1292 if (workflows == null) { 1293 workflows = new JSONArray(); 1294 } 1295 return JsonToBean.createWorkflowJobList(workflows); 1296 } 1297 else { 1298 handleError(conn); 1299 } 1300 return null; 1301 } 1302 } 1303 1304 1305 private class BundleJobInfo extends ClientCallable<BundleJob> { 1306 1307 BundleJobInfo(String jobId) { 1308 super("GET", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.JOB_SHOW_PARAM, 1309 RestConstants.JOB_SHOW_INFO)); 1310 } 1311 1312 @Override 1313 protected BundleJob call(HttpURLConnection conn) throws IOException, OozieClientException { 1314 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1315 Reader reader = new InputStreamReader(conn.getInputStream()); 1316 JSONObject json = (JSONObject) JSONValue.parse(reader); 1317 return JsonToBean.createBundleJob(json); 1318 } 1319 else { 1320 handleError(conn); 1321 } 1322 return null; 1323 } 1324 } 1325 1326 private class CoordActionInfo extends ClientCallable<CoordinatorAction> { 1327 CoordActionInfo(String actionId) { 1328 super("GET", RestConstants.JOB, notEmpty(actionId, "id"), prepareParams(RestConstants.JOB_SHOW_PARAM, 1329 RestConstants.JOB_SHOW_INFO)); 1330 } 1331 1332 @Override 1333 protected CoordinatorAction call(HttpURLConnection conn) throws IOException, OozieClientException { 1334 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1335 Reader reader = new InputStreamReader(conn.getInputStream()); 1336 JSONObject json = (JSONObject) JSONValue.parse(reader); 1337 return JsonToBean.createCoordinatorAction(json); 1338 } 1339 else { 1340 handleError(conn); 1341 } 1342 return null; 1343 } 1344 } 1345 1346 /** 1347 * Get the info of a bundle job. 1348 * 1349 * @param jobId job Id. 1350 * @return the job info. 1351 * @throws OozieClientException thrown if the job info could not be retrieved. 1352 */ 1353 public BundleJob getBundleJobInfo(String jobId) throws OozieClientException { 1354 return new BundleJobInfo(jobId).call(); 1355 } 1356 1357 /** 1358 * Get the info of a coordinator job. 1359 * 1360 * @param jobId job Id. 1361 * @return the job info. 1362 * @throws OozieClientException thrown if the job info could not be retrieved. 1363 */ 1364 public CoordinatorJob getCoordJobInfo(String jobId) throws OozieClientException { 1365 return new CoordJobInfo(jobId, null, -1, -1, "asc").call(); 1366 } 1367 1368 /** 1369 * Get the info of a coordinator job and subset actions. 1370 * 1371 * @param jobId job Id. 1372 * @param filter filter the status filter 1373 * @param start starting index in the list of actions belonging to the job 1374 * @param len number of actions to be returned 1375 * @return the job info. 1376 * @throws OozieClientException thrown if the job info could not be retrieved. 1377 */ 1378 public CoordinatorJob getCoordJobInfo(String jobId, String filter, int start, int len) 1379 throws OozieClientException { 1380 return new CoordJobInfo(jobId, filter, start, len, "asc").call(); 1381 } 1382 1383 /** 1384 * Get the info of a coordinator job and subset actions. 1385 * 1386 * @param jobId job Id. 1387 * @param filter filter the status filter 1388 * @param start starting index in the list of actions belonging to the job 1389 * @param len number of actions to be returned 1390 * @param order order to list coord actions (e.g, desc) 1391 * @return the job info. 1392 * @throws OozieClientException thrown if the job info could not be retrieved. 1393 */ 1394 public CoordinatorJob getCoordJobInfo(String jobId, String filter, int start, int len, String order) 1395 throws OozieClientException { 1396 return new CoordJobInfo(jobId, filter, start, len, order).call(); 1397 } 1398 1399 public List<WorkflowJob> getWfsForCoordAction(String coordActionId) throws OozieClientException { 1400 return new WfsForCoordAction(coordActionId).call(); 1401 } 1402 1403 /** 1404 * Get the info of a coordinator action. 1405 * 1406 * @param actionId Id. 1407 * @return the coordinator action info. 1408 * @throws OozieClientException thrown if the job info could not be retrieved. 1409 */ 1410 public CoordinatorAction getCoordActionInfo(String actionId) throws OozieClientException { 1411 return new CoordActionInfo(actionId).call(); 1412 } 1413 1414 private class JobsStatus extends ClientCallable<List<WorkflowJob>> { 1415 1416 JobsStatus(String filter, int start, int len) { 1417 super("GET", RestConstants.JOBS, "", prepareParams(RestConstants.JOBS_FILTER_PARAM, filter, 1418 RestConstants.JOBTYPE_PARAM, "wf", RestConstants.OFFSET_PARAM, Integer.toString(start), 1419 RestConstants.LEN_PARAM, Integer.toString(len))); 1420 } 1421 1422 @Override 1423 protected List<WorkflowJob> call(HttpURLConnection conn) throws IOException, OozieClientException { 1424 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1425 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1426 Reader reader = new InputStreamReader(conn.getInputStream()); 1427 JSONObject json = (JSONObject) JSONValue.parse(reader); 1428 JSONArray workflows = (JSONArray) json.get(JsonTags.WORKFLOWS_JOBS); 1429 if (workflows == null) { 1430 workflows = new JSONArray(); 1431 } 1432 return JsonToBean.createWorkflowJobList(workflows); 1433 } 1434 else { 1435 handleError(conn); 1436 } 1437 return null; 1438 } 1439 } 1440 1441 private class CoordJobsStatus extends ClientCallable<List<CoordinatorJob>> { 1442 1443 CoordJobsStatus(String filter, int start, int len) { 1444 super("GET", RestConstants.JOBS, "", prepareParams(RestConstants.JOBS_FILTER_PARAM, filter, 1445 RestConstants.JOBTYPE_PARAM, "coord", RestConstants.OFFSET_PARAM, Integer.toString(start), 1446 RestConstants.LEN_PARAM, Integer.toString(len))); 1447 } 1448 1449 @Override 1450 protected List<CoordinatorJob> call(HttpURLConnection conn) throws IOException, OozieClientException { 1451 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1452 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1453 Reader reader = new InputStreamReader(conn.getInputStream()); 1454 JSONObject json = (JSONObject) JSONValue.parse(reader); 1455 JSONArray jobs = (JSONArray) json.get(JsonTags.COORDINATOR_JOBS); 1456 if (jobs == null) { 1457 jobs = new JSONArray(); 1458 } 1459 return JsonToBean.createCoordinatorJobList(jobs); 1460 } 1461 else { 1462 handleError(conn); 1463 } 1464 return null; 1465 } 1466 } 1467 1468 private class BundleJobsStatus extends ClientCallable<List<BundleJob>> { 1469 1470 BundleJobsStatus(String filter, int start, int len) { 1471 super("GET", RestConstants.JOBS, "", prepareParams(RestConstants.JOBS_FILTER_PARAM, filter, 1472 RestConstants.JOBTYPE_PARAM, "bundle", RestConstants.OFFSET_PARAM, Integer.toString(start), 1473 RestConstants.LEN_PARAM, Integer.toString(len))); 1474 } 1475 1476 @Override 1477 protected List<BundleJob> call(HttpURLConnection conn) throws IOException, OozieClientException { 1478 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1479 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1480 Reader reader = new InputStreamReader(conn.getInputStream()); 1481 JSONObject json = (JSONObject) JSONValue.parse(reader); 1482 JSONArray jobs = (JSONArray) json.get(JsonTags.BUNDLE_JOBS); 1483 if (jobs == null) { 1484 jobs = new JSONArray(); 1485 } 1486 return JsonToBean.createBundleJobList(jobs); 1487 } 1488 else { 1489 handleError(conn); 1490 } 1491 return null; 1492 } 1493 } 1494 1495 private class BulkResponseStatus extends ClientCallable<List<BulkResponse>> { 1496 1497 BulkResponseStatus(String filter, int start, int len) { 1498 super("GET", RestConstants.JOBS, "", prepareParams(RestConstants.JOBS_BULK_PARAM, filter, 1499 RestConstants.OFFSET_PARAM, Integer.toString(start), RestConstants.LEN_PARAM, Integer.toString(len))); 1500 } 1501 1502 @Override 1503 protected List<BulkResponse> call(HttpURLConnection conn) throws IOException, OozieClientException { 1504 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1505 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1506 Reader reader = new InputStreamReader(conn.getInputStream()); 1507 JSONObject json = (JSONObject) JSONValue.parse(reader); 1508 JSONArray results = (JSONArray) json.get(JsonTags.BULK_RESPONSES); 1509 if (results == null) { 1510 results = new JSONArray(); 1511 } 1512 return JsonToBean.createBulkResponseList(results); 1513 } 1514 else { 1515 handleError(conn); 1516 } 1517 return null; 1518 } 1519 } 1520 1521 private class CoordActionsKill extends ClientCallable<List<CoordinatorAction>> { 1522 1523 CoordActionsKill(String jobId, String rangeType, String scope) { 1524 super("PUT", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.ACTION_PARAM, 1525 RestConstants.JOB_ACTION_KILL, RestConstants.JOB_COORD_RANGE_TYPE_PARAM, rangeType, 1526 RestConstants.JOB_COORD_SCOPE_PARAM, scope)); 1527 } 1528 1529 @Override 1530 protected List<CoordinatorAction> call(HttpURLConnection conn) throws IOException, OozieClientException { 1531 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1532 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1533 Reader reader = new InputStreamReader(conn.getInputStream()); 1534 JSONObject json = (JSONObject) JSONValue.parse(reader); 1535 JSONArray coordActions = (JSONArray) json.get(JsonTags.COORDINATOR_ACTIONS); 1536 return JsonToBean.createCoordinatorActionList(coordActions); 1537 } 1538 else { 1539 handleError(conn); 1540 } 1541 return null; 1542 } 1543 } 1544 1545 private class CoordIgnore extends ClientCallable<List<CoordinatorAction>> { 1546 CoordIgnore(String jobId, String rerunType, String scope) { 1547 super("PUT", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.ACTION_PARAM, 1548 RestConstants.JOB_ACTION_IGNORE, RestConstants.JOB_COORD_RANGE_TYPE_PARAM, 1549 rerunType, RestConstants.JOB_COORD_SCOPE_PARAM, scope)); 1550 } 1551 1552 @Override 1553 protected List<CoordinatorAction> call(HttpURLConnection conn) throws IOException, OozieClientException { 1554 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1555 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1556 Reader reader = new InputStreamReader(conn.getInputStream()); 1557 JSONObject json = (JSONObject) JSONValue.parse(reader); 1558 if(json != null) { 1559 JSONArray coordActions = (JSONArray) json.get(JsonTags.COORDINATOR_ACTIONS); 1560 return JsonToBean.createCoordinatorActionList(coordActions); 1561 } 1562 } 1563 else { 1564 handleError(conn); 1565 } 1566 return null; 1567 } 1568 } 1569 private class CoordRerun extends ClientCallable<List<CoordinatorAction>> { 1570 1571 CoordRerun(String jobId, String rerunType, String scope, boolean refresh, boolean noCleanup, boolean failed) { 1572 super("PUT", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.ACTION_PARAM, 1573 RestConstants.JOB_COORD_ACTION_RERUN, RestConstants.JOB_COORD_RANGE_TYPE_PARAM, rerunType, 1574 RestConstants.JOB_COORD_SCOPE_PARAM, scope, RestConstants.JOB_COORD_RERUN_REFRESH_PARAM, 1575 Boolean.toString(refresh), RestConstants.JOB_COORD_RERUN_NOCLEANUP_PARAM, Boolean 1576 .toString(noCleanup), RestConstants.JOB_COORD_RERUN_FAILED_PARAM, Boolean.toString(failed))); 1577 } 1578 1579 @Override 1580 protected List<CoordinatorAction> call(HttpURLConnection conn) throws IOException, OozieClientException { 1581 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1582 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1583 Reader reader = new InputStreamReader(conn.getInputStream()); 1584 JSONObject json = (JSONObject) JSONValue.parse(reader); 1585 JSONArray coordActions = (JSONArray) json.get(JsonTags.COORDINATOR_ACTIONS); 1586 return JsonToBean.createCoordinatorActionList(coordActions); 1587 } 1588 else { 1589 handleError(conn); 1590 } 1591 return null; 1592 } 1593 } 1594 1595 private class BundleRerun extends ClientCallable<Void> { 1596 1597 BundleRerun(String jobId, String coordScope, String dateScope, boolean refresh, boolean noCleanup) { 1598 super("PUT", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.ACTION_PARAM, 1599 RestConstants.JOB_BUNDLE_ACTION_RERUN, RestConstants.JOB_BUNDLE_RERUN_COORD_SCOPE_PARAM, 1600 coordScope, RestConstants.JOB_BUNDLE_RERUN_DATE_SCOPE_PARAM, dateScope, 1601 RestConstants.JOB_COORD_RERUN_REFRESH_PARAM, Boolean.toString(refresh), 1602 RestConstants.JOB_COORD_RERUN_NOCLEANUP_PARAM, Boolean.toString(noCleanup))); 1603 } 1604 1605 @Override 1606 protected Void call(HttpURLConnection conn) throws IOException, OozieClientException { 1607 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1608 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1609 return null; 1610 } 1611 else { 1612 handleError(conn); 1613 } 1614 return null; 1615 } 1616 } 1617 1618 /** 1619 * Rerun coordinator actions. 1620 * 1621 * @param jobId coordinator jobId 1622 * @param rerunType rerun type 'date' if -date is used, 'action-id' if -action is used 1623 * @param scope rerun scope for date or actionIds 1624 * @param refresh true if -refresh is given in command option 1625 * @param noCleanup true if -nocleanup is given in command option 1626 * @throws OozieClientException 1627 */ 1628 public List<CoordinatorAction> reRunCoord(String jobId, String rerunType, String scope, boolean refresh, 1629 boolean noCleanup) throws OozieClientException { 1630 return new CoordRerun(jobId, rerunType, scope, refresh, noCleanup, false).call(); 1631 } 1632 1633 /** 1634 * Rerun coordinator actions with failed option. 1635 * 1636 * @param jobId coordinator jobId 1637 * @param rerunType rerun type 'date' if -date is used, 'action-id' if -action is used 1638 * @param scope rerun scope for date or actionIds 1639 * @param refresh true if -refresh is given in command option 1640 * @param noCleanup true if -nocleanup is given in command option 1641 * @param failed true if -failed is given in command option 1642 * @throws OozieClientException 1643 */ 1644 public List<CoordinatorAction> reRunCoord(String jobId, String rerunType, String scope, boolean refresh, 1645 boolean noCleanup, boolean failed) throws OozieClientException { 1646 return new CoordRerun(jobId, rerunType, scope, refresh, noCleanup, failed).call(); 1647 } 1648 1649 /** 1650 * Rerun bundle coordinators. 1651 * 1652 * @param jobId bundle jobId 1653 * @param coordScope rerun scope for coordinator jobs 1654 * @param dateScope rerun scope for date 1655 * @param refresh true if -refresh is given in command option 1656 * @param noCleanup true if -nocleanup is given in command option 1657 * @throws OozieClientException 1658 */ 1659 public Void reRunBundle(String jobId, String coordScope, String dateScope, boolean refresh, boolean noCleanup) 1660 throws OozieClientException { 1661 return new BundleRerun(jobId, coordScope, dateScope, refresh, noCleanup).call(); 1662 } 1663 1664 /** 1665 * Return the info of the workflow jobs that match the filter. 1666 * 1667 * @param filter job filter. Refer to the {@link OozieClient} for the filter syntax. 1668 * @param start jobs offset, base 1. 1669 * @param len number of jobs to return. 1670 * @return a list with the workflow jobs info, without node details. 1671 * @throws OozieClientException thrown if the jobs info could not be retrieved. 1672 */ 1673 public List<WorkflowJob> getJobsInfo(String filter, int start, int len) throws OozieClientException { 1674 return new JobsStatus(filter, start, len).call(); 1675 } 1676 1677 /** 1678 * Return the info of the workflow jobs that match the filter. 1679 * <p/> 1680 * It returns the first 100 jobs that match the filter. 1681 * 1682 * @param filter job filter. Refer to the {@link OozieClient} for the filter syntax. 1683 * @return a list with the workflow jobs info, without node details. 1684 * @throws OozieClientException thrown if the jobs info could not be retrieved. 1685 */ 1686 public List<WorkflowJob> getJobsInfo(String filter) throws OozieClientException { 1687 return getJobsInfo(filter, 1, 50); 1688 } 1689 1690 /** 1691 * Print sla info about coordinator and workflow jobs and actions. 1692 * 1693 * @param start starting offset 1694 * @param len number of results 1695 * @throws OozieClientException 1696 */ 1697 public void getSlaInfo(int start, int len, String filter) throws OozieClientException { 1698 new SlaInfo(start, len, filter).call(); 1699 } 1700 1701 private class SlaInfo extends ClientCallable<Void> { 1702 1703 SlaInfo(int start, int len, String filter) { 1704 super("GET", WS_PROTOCOL_VERSION_1, RestConstants.SLA, "", prepareParams(RestConstants.SLA_GT_SEQUENCE_ID, 1705 Integer.toString(start), RestConstants.MAX_EVENTS, Integer.toString(len), 1706 RestConstants.JOBS_FILTER_PARAM, filter)); 1707 } 1708 1709 @Override 1710 protected Void call(HttpURLConnection conn) throws IOException, OozieClientException { 1711 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1712 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1713 BufferedReader br = new BufferedReader(new InputStreamReader(conn.getInputStream())); 1714 String line = null; 1715 while ((line = br.readLine()) != null) { 1716 System.out.println(line); 1717 } 1718 } 1719 else { 1720 handleError(conn); 1721 } 1722 return null; 1723 } 1724 } 1725 1726 private class JobIdAction extends ClientCallable<String> { 1727 1728 JobIdAction(String externalId) { 1729 super("GET", RestConstants.JOBS, "", prepareParams(RestConstants.JOBTYPE_PARAM, "wf", 1730 RestConstants.JOBS_EXTERNAL_ID_PARAM, externalId)); 1731 } 1732 1733 @Override 1734 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 1735 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1736 Reader reader = new InputStreamReader(conn.getInputStream()); 1737 JSONObject json = (JSONObject) JSONValue.parse(reader); 1738 return (String) json.get(JsonTags.JOB_ID); 1739 } 1740 else { 1741 handleError(conn); 1742 } 1743 return null; 1744 } 1745 } 1746 1747 /** 1748 * Return the workflow job Id for an external Id. 1749 * <p/> 1750 * The external Id must have provided at job creation time. 1751 * 1752 * @param externalId external Id given at job creation time. 1753 * @return the workflow job Id for an external Id, <code>null</code> if none. 1754 * @throws OozieClientException thrown if the operation could not be done. 1755 */ 1756 public String getJobId(String externalId) throws OozieClientException { 1757 return new JobIdAction(externalId).call(); 1758 } 1759 1760 private class SetSystemMode extends ClientCallable<Void> { 1761 1762 public SetSystemMode(SYSTEM_MODE status) { 1763 super("PUT", RestConstants.ADMIN, RestConstants.ADMIN_STATUS_RESOURCE, prepareParams( 1764 RestConstants.ADMIN_SYSTEM_MODE_PARAM, status + "")); 1765 } 1766 1767 @Override 1768 public Void call(HttpURLConnection conn) throws IOException, OozieClientException { 1769 if (conn.getResponseCode() != HttpURLConnection.HTTP_OK) { 1770 handleError(conn); 1771 } 1772 return null; 1773 } 1774 } 1775 1776 /** 1777 * Enable or disable safe mode. Used by OozieCLI. In safe mode, Oozie would not accept any commands except status 1778 * command to change and view the safe mode status. 1779 * 1780 * @param status true to enable safe mode, false to disable safe mode. 1781 * @throws OozieClientException if it fails to set the safe mode status. 1782 */ 1783 public void setSystemMode(SYSTEM_MODE status) throws OozieClientException { 1784 new SetSystemMode(status).call(); 1785 } 1786 1787 private class GetSystemMode extends ClientCallable<SYSTEM_MODE> { 1788 1789 GetSystemMode() { 1790 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_STATUS_RESOURCE, prepareParams()); 1791 } 1792 1793 @Override 1794 protected SYSTEM_MODE call(HttpURLConnection conn) throws IOException, OozieClientException { 1795 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1796 Reader reader = new InputStreamReader(conn.getInputStream()); 1797 JSONObject json = (JSONObject) JSONValue.parse(reader); 1798 return SYSTEM_MODE.valueOf((String) json.get(JsonTags.OOZIE_SYSTEM_MODE)); 1799 } 1800 else { 1801 handleError(conn); 1802 } 1803 return SYSTEM_MODE.NORMAL; 1804 } 1805 } 1806 1807 /** 1808 * Returns if Oozie is in safe mode or not. 1809 * 1810 * @return true if safe mode is ON<br> 1811 * false if safe mode is OFF 1812 * @throws OozieClientException throw if it could not obtain the safe mode status. 1813 */ 1814 /* 1815 * public boolean isInSafeMode() throws OozieClientException { return new GetSafeMode().call(); } 1816 */ 1817 public SYSTEM_MODE getSystemMode() throws OozieClientException { 1818 return new GetSystemMode().call(); 1819 } 1820 1821 public String updateShareLib() throws OozieClientException { 1822 return new UpdateSharelib().call(); 1823 } 1824 1825 public String listShareLib(String sharelibKey) throws OozieClientException { 1826 return new ListShareLib(sharelibKey).call(); 1827 } 1828 1829 private class GetBuildVersion extends ClientCallable<String> { 1830 1831 GetBuildVersion() { 1832 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_BUILD_VERSION_RESOURCE, prepareParams()); 1833 } 1834 1835 @Override 1836 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 1837 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1838 Reader reader = new InputStreamReader(conn.getInputStream()); 1839 JSONObject json = (JSONObject) JSONValue.parse(reader); 1840 return (String) json.get(JsonTags.BUILD_VERSION); 1841 } 1842 else { 1843 handleError(conn); 1844 } 1845 return null; 1846 } 1847 } 1848 1849 private class ValidateXML extends ClientCallable<String> { 1850 1851 String file = null; 1852 1853 ValidateXML(String file, String user) { 1854 super("POST", RestConstants.VALIDATE, "", 1855 prepareParams(RestConstants.FILE_PARAM, file, RestConstants.USER_PARAM, user)); 1856 this.file = file; 1857 } 1858 1859 @Override 1860 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 1861 conn.setRequestProperty("content-type", RestConstants.XML_CONTENT_TYPE); 1862 if (file.startsWith("/")) { 1863 FileInputStream fi = new FileInputStream(new File(file)); 1864 byte[] buffer = new byte[1024]; 1865 int n = 0; 1866 while (-1 != (n = fi.read(buffer))) { 1867 conn.getOutputStream().write(buffer, 0, n); 1868 } 1869 } 1870 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1871 Reader reader = new InputStreamReader(conn.getInputStream()); 1872 JSONObject json = (JSONObject) JSONValue.parse(reader); 1873 return (String) json.get(JsonTags.VALIDATE); 1874 } 1875 else if ((conn.getResponseCode() == HttpURLConnection.HTTP_NOT_FOUND)) { 1876 return null; 1877 } 1878 else { 1879 handleError(conn); 1880 } 1881 return null; 1882 } 1883 } 1884 1885 1886 private class UpdateSharelib extends ClientCallable<String> { 1887 1888 UpdateSharelib() { 1889 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_UPDATE_SHARELIB, prepareParams( 1890 RestConstants.ALL_SERVER_REQUEST, "true")); 1891 } 1892 1893 @Override 1894 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 1895 StringBuffer bf = new StringBuffer(); 1896 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1897 Reader reader = new InputStreamReader(conn.getInputStream()); 1898 Object sharelib = (Object) JSONValue.parse(reader); 1899 bf.append("[ShareLib update status]").append(System.getProperty("line.separator")); 1900 if (sharelib instanceof JSONArray) { 1901 for (Object o : ((JSONArray) sharelib)) { 1902 JSONObject obj = (JSONObject) ((JSONObject) o).get(JsonTags.SHARELIB_LIB_UPDATE); 1903 for (Object key : obj.keySet()) { 1904 bf.append("\t").append(key).append(" = ").append(obj.get(key)) 1905 .append(System.getProperty("line.separator")); 1906 } 1907 bf.append(System.getProperty("line.separator")); 1908 } 1909 } 1910 else{ 1911 JSONObject obj = (JSONObject) ((JSONObject) sharelib).get(JsonTags.SHARELIB_LIB_UPDATE); 1912 for (Object key : obj.keySet()) { 1913 bf.append("\t").append(key).append(" = ").append(obj.get(key)) 1914 .append(System.getProperty("line.separator")); 1915 } 1916 bf.append(System.getProperty("line.separator")); 1917 } 1918 return bf.toString(); 1919 } 1920 else { 1921 handleError(conn); 1922 } 1923 return null; 1924 } 1925 } 1926 1927 private class ListShareLib extends ClientCallable<String> { 1928 1929 ListShareLib(String sharelibKey) { 1930 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_LIST_SHARELIB, prepareParams( 1931 RestConstants.SHARE_LIB_REQUEST_KEY, sharelibKey)); 1932 } 1933 1934 @Override 1935 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 1936 1937 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 1938 StringBuffer bf = new StringBuffer(); 1939 Reader reader = new InputStreamReader(conn.getInputStream()); 1940 JSONObject json = (JSONObject) JSONValue.parse(reader); 1941 Object sharelib = json.get(JsonTags.SHARELIB_LIB); 1942 bf.append("[Available ShareLib]").append(System.getProperty("line.separator")); 1943 if (sharelib instanceof JSONArray) { 1944 for (Object o : ((JSONArray) sharelib)) { 1945 JSONObject obj = (JSONObject) o; 1946 bf.append(obj.get(JsonTags.SHARELIB_LIB_NAME)) 1947 .append(System.getProperty("line.separator")); 1948 if (obj.get(JsonTags.SHARELIB_LIB_FILES) != null) { 1949 for (Object file : ((JSONArray) obj.get(JsonTags.SHARELIB_LIB_FILES))) { 1950 bf.append("\t").append(file).append(System.getProperty("line.separator")); 1951 } 1952 } 1953 } 1954 return bf.toString(); 1955 } 1956 } 1957 else { 1958 handleError(conn); 1959 } 1960 return null; 1961 } 1962 1963 } 1964 1965 /** 1966 * Return the Oozie server build version. 1967 * 1968 * @return the Oozie server build version. 1969 * @throws OozieClientException throw if it the server build version could not be retrieved. 1970 */ 1971 public String getServerBuildVersion() throws OozieClientException { 1972 return new GetBuildVersion().call(); 1973 } 1974 1975 /** 1976 * Return the Oozie client build version. 1977 * 1978 * @return the Oozie client build version. 1979 */ 1980 public String getClientBuildVersion() { 1981 return BuildInfo.getBuildInfo().getProperty(BuildInfo.BUILD_VERSION); 1982 } 1983 1984 /** 1985 * Return the workflow application is valid. 1986 * 1987 * @param file local file or hdfs file. 1988 * @return the workflow application is valid. 1989 * @throws OozieClientException throw if it the workflow application's validation could not be retrieved. 1990 */ 1991 public String validateXML(String file) throws OozieClientException { 1992 String fileName = file; 1993 if (file.startsWith("file://")) { 1994 fileName = file.substring(7, file.length()); 1995 } 1996 if (!fileName.contains("://")) { 1997 File f = new File(fileName); 1998 if (!f.isFile()) { 1999 throw new OozieClientException("File error", "File does not exist : " + f.getAbsolutePath()); 2000 } 2001 fileName = f.getAbsolutePath(); 2002 } 2003 String user = USER_NAME_TL.get(); 2004 if (user == null) { 2005 user = System.getProperty("user.name"); 2006 } 2007 return new ValidateXML(fileName, user).call(); 2008 } 2009 2010 /** 2011 * Return the info of the coordinator jobs that match the filter. 2012 * 2013 * @param filter job filter. Refer to the {@link OozieClient} for the filter syntax. 2014 * @param start jobs offset, base 1. 2015 * @param len number of jobs to return. 2016 * @return a list with the coordinator jobs info 2017 * @throws OozieClientException thrown if the jobs info could not be retrieved. 2018 */ 2019 public List<CoordinatorJob> getCoordJobsInfo(String filter, int start, int len) throws OozieClientException { 2020 return new CoordJobsStatus(filter, start, len).call(); 2021 } 2022 2023 /** 2024 * Return the info of the bundle jobs that match the filter. 2025 * 2026 * @param filter job filter. Refer to the {@link OozieClient} for the filter syntax. 2027 * @param start jobs offset, base 1. 2028 * @param len number of jobs to return. 2029 * @return a list with the bundle jobs info 2030 * @throws OozieClientException thrown if the jobs info could not be retrieved. 2031 */ 2032 public List<BundleJob> getBundleJobsInfo(String filter, int start, int len) throws OozieClientException { 2033 return new BundleJobsStatus(filter, start, len).call(); 2034 } 2035 2036 public List<BulkResponse> getBulkInfo(String filter, int start, int len) throws OozieClientException { 2037 return new BulkResponseStatus(filter, start, len).call(); 2038 } 2039 2040 /** 2041 * Poll a job (Workflow Job ID, Coordinator Job ID, Coordinator Action ID, or Bundle Job ID) and return when it has reached a 2042 * terminal state. 2043 * (i.e. FAILED, KILLED, SUCCEEDED) 2044 * 2045 * @param id The Job ID 2046 * @param timeout timeout in minutes (negative values indicate no timeout) 2047 * @param interval polling interval in minutes (must be positive) 2048 * @param verbose if true, the current status will be printed out at each poll; if false, no output 2049 * @throws OozieClientException thrown if the job's status could not be retrieved 2050 */ 2051 public void pollJob(String id, int timeout, int interval, boolean verbose) throws OozieClientException { 2052 notEmpty("id", id); 2053 if (interval < 1) { 2054 throw new IllegalArgumentException("interval must be a positive integer"); 2055 } 2056 boolean noTimeout = (timeout < 1); 2057 long endTime = System.currentTimeMillis() + timeout * 60 * 1000; 2058 interval *= 60 * 1000; 2059 2060 final Set<String> completedStatuses; 2061 if (id.endsWith("-W")) { 2062 completedStatuses = COMPLETED_WF_STATUSES; 2063 } else if (id.endsWith("-C")) { 2064 completedStatuses = COMPLETED_COORD_AND_BUNDLE_STATUSES; 2065 } else if (id.endsWith("-B")) { 2066 completedStatuses = COMPLETED_COORD_AND_BUNDLE_STATUSES; 2067 } else if (id.contains("-C@")) { 2068 completedStatuses = COMPLETED_COORD_ACTION_STATUSES; 2069 } else { 2070 throw new IllegalArgumentException("invalid job type"); 2071 } 2072 2073 String status = getStatus(id); 2074 if (verbose) { 2075 System.out.println(status); 2076 } 2077 while(!completedStatuses.contains(status) && (noTimeout || System.currentTimeMillis() <= endTime)) { 2078 try { 2079 Thread.sleep(interval); 2080 } catch (InterruptedException ie) { 2081 // ignore 2082 } 2083 status = getStatus(id); 2084 if (verbose) { 2085 System.out.println(status); 2086 } 2087 } 2088 } 2089 2090 /** 2091 * Gets the status for a particular job (Workflow Job ID, Coordinator Job ID, Coordinator Action ID, or Bundle Job ID). 2092 * 2093 * @param jobId given jobId 2094 * @return the status 2095 * @throws OozieClientException 2096 */ 2097 public String getStatus(String jobId) throws OozieClientException { 2098 return new Status(jobId).call(); 2099 } 2100 2101 private class Status extends ClientCallable<String> { 2102 2103 Status(String jobId) { 2104 super("GET", RestConstants.JOB, notEmpty(jobId, "jobId"), prepareParams(RestConstants.JOB_SHOW_PARAM, 2105 RestConstants.JOB_SHOW_STATUS)); 2106 } 2107 2108 @Override 2109 protected String call(HttpURLConnection conn) throws IOException, OozieClientException { 2110 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 2111 Reader reader = new InputStreamReader(conn.getInputStream()); 2112 JSONObject json = (JSONObject) JSONValue.parse(reader); 2113 return (String) json.get(JsonTags.STATUS); 2114 } 2115 else { 2116 handleError(conn); 2117 } 2118 return null; 2119 } 2120 } 2121 2122 private class GetQueueDump extends ClientCallable<List<String>> { 2123 GetQueueDump() { 2124 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_QUEUE_DUMP_RESOURCE, prepareParams()); 2125 } 2126 2127 @Override 2128 protected List<String> call(HttpURLConnection conn) throws IOException, OozieClientException { 2129 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 2130 Reader reader = new InputStreamReader(conn.getInputStream()); 2131 JSONObject json = (JSONObject) JSONValue.parse(reader); 2132 JSONArray queueDumpArray = (JSONArray) json.get(JsonTags.QUEUE_DUMP); 2133 2134 List<String> list = new ArrayList<String>(); 2135 list.add("[Server Queue Dump]:"); 2136 for (Object o : queueDumpArray) { 2137 JSONObject entry = (JSONObject) o; 2138 if (entry.get(JsonTags.CALLABLE_DUMP) != null) { 2139 String value = (String) entry.get(JsonTags.CALLABLE_DUMP); 2140 list.add(value); 2141 } 2142 } 2143 if (queueDumpArray.size() == 0) { 2144 list.add("Queue dump is null!"); 2145 } 2146 2147 list.add("******************************************"); 2148 list.add("[Server Uniqueness Map Dump]:"); 2149 2150 JSONArray uniqueDumpArray = (JSONArray) json.get(JsonTags.UNIQUE_MAP_DUMP); 2151 for (Object o : uniqueDumpArray) { 2152 JSONObject entry = (JSONObject) o; 2153 if (entry.get(JsonTags.UNIQUE_ENTRY_DUMP) != null) { 2154 String value = (String) entry.get(JsonTags.UNIQUE_ENTRY_DUMP); 2155 list.add(value); 2156 } 2157 } 2158 if (uniqueDumpArray.size() == 0) { 2159 list.add("Uniqueness dump is null!"); 2160 } 2161 return list; 2162 } 2163 else { 2164 handleError(conn); 2165 } 2166 return null; 2167 } 2168 } 2169 2170 /** 2171 * Return the Oozie queue's commands' dump 2172 * 2173 * @return the list of strings of callable identification in queue 2174 * @throws OozieClientException throw if it the queue dump could not be retrieved. 2175 */ 2176 public List<String> getQueueDump() throws OozieClientException { 2177 return new GetQueueDump().call(); 2178 } 2179 2180 private class GetAvailableOozieServers extends MapClientCallable { 2181 2182 GetAvailableOozieServers() { 2183 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_AVAILABLE_OOZIE_SERVERS_RESOURCE, prepareParams()); 2184 } 2185 } 2186 2187 /** 2188 * Return the list of available Oozie servers. 2189 * 2190 * @return the list of available Oozie servers. 2191 * @throws OozieClientException throw if it the list of available Oozie servers could not be retrieved. 2192 */ 2193 public Map<String, String> getAvailableOozieServers() throws OozieClientException { 2194 return new GetAvailableOozieServers().call(); 2195 } 2196 2197 private class GetServerConfiguration extends MapClientCallable { 2198 2199 GetServerConfiguration() { 2200 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_CONFIG_RESOURCE, prepareParams()); 2201 } 2202 } 2203 2204 /** 2205 * Return the Oozie system configuration. 2206 * 2207 * @return the Oozie system configuration. 2208 * @throws OozieClientException throw if the system configuration could not be retrieved. 2209 */ 2210 public Map<String, String> getServerConfiguration() throws OozieClientException { 2211 return new GetServerConfiguration().call(); 2212 } 2213 2214 private class GetJavaSystemProperties extends MapClientCallable { 2215 2216 GetJavaSystemProperties() { 2217 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_JAVA_SYS_PROPS_RESOURCE, prepareParams()); 2218 } 2219 } 2220 2221 /** 2222 * Return the Oozie Java system properties. 2223 * 2224 * @return the Oozie Java system properties. 2225 * @throws OozieClientException throw if the system properties could not be retrieved. 2226 */ 2227 public Map<String, String> getJavaSystemProperties() throws OozieClientException { 2228 return new GetJavaSystemProperties().call(); 2229 } 2230 2231 private class GetOSEnv extends MapClientCallable { 2232 2233 GetOSEnv() { 2234 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_OS_ENV_RESOURCE, prepareParams()); 2235 } 2236 } 2237 2238 /** 2239 * Return the Oozie system OS environment. 2240 * 2241 * @return the Oozie system OS environment. 2242 * @throws OozieClientException throw if the system OS environment could not be retrieved. 2243 */ 2244 public Map<String, String> getOSEnv() throws OozieClientException { 2245 return new GetOSEnv().call(); 2246 } 2247 2248 private class GetMetrics extends ClientCallable<Metrics> { 2249 2250 GetMetrics() { 2251 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_METRICS_RESOURCE, prepareParams()); 2252 } 2253 2254 @Override 2255 protected Metrics call(HttpURLConnection conn) throws IOException, OozieClientException { 2256 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 2257 Reader reader = new InputStreamReader(conn.getInputStream()); 2258 JSONObject json = (JSONObject) JSONValue.parse(reader); 2259 Metrics metrics = new Metrics(json); 2260 return metrics; 2261 } 2262 else if ((conn.getResponseCode() == HttpURLConnection.HTTP_UNAVAILABLE)) { 2263 // Use Instrumentation endpoint 2264 return null; 2265 } 2266 else { 2267 handleError(conn); 2268 } 2269 return null; 2270 } 2271 } 2272 2273 public class Metrics { 2274 private Map<String, Long> counters; 2275 private Map<String, Object> gauges; 2276 private Map<String, Timer> timers; 2277 private Map<String, Histogram> histograms; 2278 2279 @SuppressWarnings("unchecked") 2280 public Metrics(JSONObject json) { 2281 JSONObject jCounters = (JSONObject) json.get("counters"); 2282 counters = new HashMap<String, Long>(jCounters.size()); 2283 for (Object entO : jCounters.entrySet()) { 2284 Entry<String, JSONObject> ent = (Entry<String, JSONObject>) entO; 2285 counters.put(ent.getKey(), (Long)ent.getValue().get("count")); 2286 } 2287 2288 JSONObject jGuages = (JSONObject) json.get("gauges"); 2289 gauges = new HashMap<String, Object>(jGuages.size()); 2290 for (Object entO : jGuages.entrySet()) { 2291 Entry<String, JSONObject> ent = (Entry<String, JSONObject>) entO; 2292 gauges.put(ent.getKey(), ent.getValue().get("value")); 2293 } 2294 2295 JSONObject jTimers = (JSONObject) json.get("timers"); 2296 timers = new HashMap<String, Timer>(jTimers.size()); 2297 for (Object entO : jTimers.entrySet()) { 2298 Entry<String, JSONObject> ent = (Entry<String, JSONObject>) entO; 2299 timers.put(ent.getKey(), new Timer(ent.getValue())); 2300 } 2301 2302 JSONObject jHistograms = (JSONObject) json.get("histograms"); 2303 histograms = new HashMap<String, Histogram>(jHistograms.size()); 2304 for (Object entO : jHistograms.entrySet()) { 2305 Entry<String, JSONObject> ent = (Entry<String, JSONObject>) entO; 2306 histograms.put(ent.getKey(), new Histogram(ent.getValue())); 2307 } 2308 } 2309 2310 public Map<String, Long> getCounters() { 2311 return counters; 2312 } 2313 2314 public Map<String, Object> getGauges() { 2315 return gauges; 2316 } 2317 2318 public Map<String, Timer> getTimers() { 2319 return timers; 2320 } 2321 2322 public Map<String, Histogram> getHistograms() { 2323 return histograms; 2324 } 2325 2326 public class Timer extends Histogram { 2327 private double m15Rate; 2328 private double m5Rate; 2329 private double m1Rate; 2330 private double meanRate; 2331 private String durationUnits; 2332 private String rateUnits; 2333 2334 public Timer(JSONObject json) { 2335 super(json); 2336 m15Rate = Double.valueOf(json.get("m15_rate").toString()); 2337 m5Rate = Double.valueOf(json.get("m5_rate").toString()); 2338 m1Rate = Double.valueOf(json.get("m1_rate").toString()); 2339 meanRate = Double.valueOf(json.get("mean_rate").toString()); 2340 durationUnits = json.get("duration_units").toString(); 2341 rateUnits = json.get("rate_units").toString(); 2342 } 2343 2344 public double get15MinuteRate() { 2345 return m15Rate; 2346 } 2347 2348 public double get5MinuteRate() { 2349 return m5Rate; 2350 } 2351 2352 public double get1MinuteRate() { 2353 return m1Rate; 2354 } 2355 2356 public double getMeanRate() { 2357 return meanRate; 2358 } 2359 2360 public String getDurationUnits() { 2361 return durationUnits; 2362 } 2363 2364 public String getRateUnits() { 2365 return rateUnits; 2366 } 2367 2368 @Override 2369 public String toString() { 2370 StringBuilder sb = new StringBuilder(super.toString()); 2371 sb.append("\n\t15 minute rate : ").append(m15Rate); 2372 sb.append("\n\t5 minute rate : ").append(m5Rate); 2373 sb.append("\n\t1 minute rate : ").append(m15Rate); 2374 sb.append("\n\tmean rate : ").append(meanRate); 2375 sb.append("\n\tduration units : ").append(durationUnits); 2376 sb.append("\n\trate units : ").append(rateUnits); 2377 return sb.toString(); 2378 } 2379 } 2380 2381 public class Histogram { 2382 private double p999; 2383 private double p99; 2384 private double p98; 2385 private double p95; 2386 private double p75; 2387 private double p50; 2388 private double mean; 2389 private double min; 2390 private double max; 2391 private double stdDev; 2392 private long count; 2393 2394 public Histogram(JSONObject json) { 2395 p999 = Double.valueOf(json.get("p999").toString()); 2396 p99 = Double.valueOf(json.get("p99").toString()); 2397 p98 = Double.valueOf(json.get("p98").toString()); 2398 p95 = Double.valueOf(json.get("p95").toString()); 2399 p75 = Double.valueOf(json.get("p75").toString()); 2400 p50 = Double.valueOf(json.get("p50").toString()); 2401 mean = Double.valueOf(json.get("mean").toString()); 2402 min = Double.valueOf(json.get("min").toString()); 2403 max = Double.valueOf(json.get("max").toString()); 2404 stdDev = Double.valueOf(json.get("stddev").toString()); 2405 count = Long.valueOf(json.get("count").toString()); 2406 } 2407 2408 public double get999thPercentile() { 2409 return p999; 2410 } 2411 2412 public double get99thPercentile() { 2413 return p99; 2414 } 2415 2416 public double get98thPercentile() { 2417 return p98; 2418 } 2419 2420 public double get95thPercentile() { 2421 return p95; 2422 } 2423 2424 public double get75thPercentile() { 2425 return p75; 2426 } 2427 2428 public double get50thPercentile() { 2429 return p50; 2430 } 2431 2432 public double getMean() { 2433 return mean; 2434 } 2435 2436 public double getMin() { 2437 return min; 2438 } 2439 2440 public double getMax() { 2441 return max; 2442 } 2443 2444 public double getStandardDeviation() { 2445 return stdDev; 2446 } 2447 2448 public long getCount() { 2449 return count; 2450 } 2451 2452 @Override 2453 public String toString() { 2454 StringBuilder sb = new StringBuilder(); 2455 sb.append("\t999th percentile : ").append(p999); 2456 sb.append("\n\t99th percentile : ").append(p99); 2457 sb.append("\n\t98th percentile : ").append(p98); 2458 sb.append("\n\t95th percentile : ").append(p95); 2459 sb.append("\n\t75th percentile : ").append(p75); 2460 sb.append("\n\t50th percentile : ").append(p50); 2461 sb.append("\n\tmean : ").append(mean); 2462 sb.append("\n\tmax : ").append(max); 2463 sb.append("\n\tmin : ").append(min); 2464 sb.append("\n\tcount : ").append(count); 2465 sb.append("\n\tstandard deviation : ").append(stdDev); 2466 return sb.toString(); 2467 } 2468 } 2469 } 2470 2471 /** 2472 * Return the Oozie metrics. If null is returned, then try {@link #getInstrumentation()}. 2473 * 2474 * @return the Oozie metrics or null. 2475 * @throws OozieClientException throw if the metrics could not be retrieved. 2476 */ 2477 public Metrics getMetrics() throws OozieClientException { 2478 return new GetMetrics().call(); 2479 } 2480 2481 private class GetInstrumentation extends ClientCallable<Instrumentation> { 2482 2483 GetInstrumentation() { 2484 super("GET", RestConstants.ADMIN, RestConstants.ADMIN_INSTRUMENTATION_RESOURCE, prepareParams()); 2485 } 2486 2487 @Override 2488 protected Instrumentation call(HttpURLConnection conn) throws IOException, OozieClientException { 2489 if ((conn.getResponseCode() == HttpURLConnection.HTTP_OK)) { 2490 Reader reader = new InputStreamReader(conn.getInputStream()); 2491 JSONObject json = (JSONObject) JSONValue.parse(reader); 2492 Instrumentation instrumentation = new Instrumentation(json); 2493 return instrumentation; 2494 } 2495 else if ((conn.getResponseCode() == HttpURLConnection.HTTP_UNAVAILABLE)) { 2496 // Use Metrics endpoint 2497 return null; 2498 } 2499 else { 2500 handleError(conn); 2501 } 2502 return null; 2503 } 2504 } 2505 2506 public class Instrumentation { 2507 private Map<String, Long> counters; 2508 private Map<String, Object> variables; 2509 private Map<String, Double> samplers; 2510 private Map<String, Timer> timers; 2511 2512 public Instrumentation(JSONObject json) { 2513 JSONArray jCounters = (JSONArray) json.get("counters"); 2514 counters = new HashMap<String, Long>(jCounters.size()); 2515 for (Object groupO : jCounters) { 2516 JSONObject group = (JSONObject) groupO; 2517 String groupName = group.get("group").toString() + "."; 2518 JSONArray data = (JSONArray) group.get("data"); 2519 for (Object datO : data) { 2520 JSONObject dat = (JSONObject) datO; 2521 counters.put(groupName + dat.get("name").toString(), Long.valueOf(dat.get("value").toString())); 2522 } 2523 } 2524 2525 JSONArray jVariables = (JSONArray) json.get("variables"); 2526 variables = new HashMap<String, Object>(jVariables.size()); 2527 for (Object groupO : jVariables) { 2528 JSONObject group = (JSONObject) groupO; 2529 String groupName = group.get("group").toString() + "."; 2530 JSONArray data = (JSONArray) group.get("data"); 2531 for (Object datO : data) { 2532 JSONObject dat = (JSONObject) datO; 2533 variables.put(groupName + dat.get("name").toString(), dat.get("value")); 2534 } 2535 } 2536 2537 JSONArray jSamplers = (JSONArray) json.get("samplers"); 2538 samplers = new HashMap<String, Double>(jSamplers.size()); 2539 for (Object groupO : jSamplers) { 2540 JSONObject group = (JSONObject) groupO; 2541 String groupName = group.get("group").toString() + "."; 2542 JSONArray data = (JSONArray) group.get("data"); 2543 for (Object datO : data) { 2544 JSONObject dat = (JSONObject) datO; 2545 samplers.put(groupName + dat.get("name").toString(), Double.valueOf(dat.get("value").toString())); 2546 } 2547 } 2548 2549 JSONArray jTimers = (JSONArray) json.get("timers"); 2550 timers = new HashMap<String, Timer>(jTimers.size()); 2551 for (Object groupO : jTimers) { 2552 JSONObject group = (JSONObject) groupO; 2553 String groupName = group.get("group").toString() + "."; 2554 JSONArray data = (JSONArray) group.get("data"); 2555 for (Object datO : data) { 2556 JSONObject dat = (JSONObject) datO; 2557 timers.put(groupName + dat.get("name").toString(), new Timer(dat)); 2558 } 2559 } 2560 } 2561 2562 public class Timer { 2563 private double ownTimeStdDev; 2564 private long ownTimeAvg; 2565 private long ownMaxTime; 2566 private long ownMinTime; 2567 private double totalTimeStdDev; 2568 private long totalTimeAvg; 2569 private long totalMaxTime; 2570 private long totalMinTime; 2571 private long ticks; 2572 2573 public Timer(JSONObject json) { 2574 ownTimeStdDev = Double.valueOf(json.get("ownTimeStdDev").toString()); 2575 ownTimeAvg = Long.valueOf(json.get("ownTimeAvg").toString()); 2576 ownMaxTime = Long.valueOf(json.get("ownMaxTime").toString()); 2577 ownMinTime = Long.valueOf(json.get("ownMinTime").toString()); 2578 totalTimeStdDev = Double.valueOf(json.get("totalTimeStdDev").toString()); 2579 totalTimeAvg = Long.valueOf(json.get("totalTimeAvg").toString()); 2580 totalMaxTime = Long.valueOf(json.get("totalMaxTime").toString()); 2581 totalMinTime = Long.valueOf(json.get("totalMinTime").toString()); 2582 ticks = Long.valueOf(json.get("ticks").toString()); 2583 } 2584 2585 public double getOwnTimeStandardDeviation() { 2586 return ownTimeStdDev; 2587 } 2588 2589 public long getOwnTimeAverage() { 2590 return ownTimeAvg; 2591 } 2592 2593 public long getOwnMaxTime() { 2594 return ownMaxTime; 2595 } 2596 2597 public long getOwnMinTime() { 2598 return ownMinTime; 2599 } 2600 2601 public double getTotalTimeStandardDeviation() { 2602 return totalTimeStdDev; 2603 } 2604 2605 public long getTotalTimeAverage() { 2606 return totalTimeAvg; 2607 } 2608 2609 public long getTotalMaxTime() { 2610 return totalMaxTime; 2611 } 2612 2613 public long getTotalMinTime() { 2614 return totalMinTime; 2615 } 2616 2617 public long getTicks() { 2618 return ticks; 2619 } 2620 2621 @Override 2622 public String toString() { 2623 StringBuilder sb = new StringBuilder(); 2624 sb.append("\town time standard deviation : ").append(ownTimeStdDev); 2625 sb.append("\n\town average time : ").append(ownTimeAvg); 2626 sb.append("\n\town max time : ").append(ownMaxTime); 2627 sb.append("\n\town min time : ").append(ownMinTime); 2628 sb.append("\n\ttotal time standard deviation : ").append(totalTimeStdDev); 2629 sb.append("\n\ttotal average time : ").append(totalTimeAvg); 2630 sb.append("\n\ttotal max time : ").append(totalMaxTime); 2631 sb.append("\n\ttotal min time : ").append(totalMinTime); 2632 sb.append("\n\tticks : ").append(ticks); 2633 return sb.toString(); 2634 } 2635 } 2636 2637 public Map<String, Long> getCounters() { 2638 return counters; 2639 } 2640 2641 public Map<String, Object> getVariables() { 2642 return variables; 2643 } 2644 2645 public Map<String, Double> getSamplers() { 2646 return samplers; 2647 } 2648 2649 public Map<String, Timer> getTimers() { 2650 return timers; 2651 } 2652 } 2653 2654 /** 2655 * Return the Oozie instrumentation. If null is returned, then try {@link #getMetrics()}. 2656 * 2657 * @return the Oozie intstrumentation or null. 2658 * @throws OozieClientException throw if the intstrumentation could not be retrieved. 2659 */ 2660 public Instrumentation getInstrumentation() throws OozieClientException { 2661 return new GetInstrumentation().call(); 2662 } 2663 2664 /** 2665 * Check if the string is not null or not empty. 2666 * 2667 * @param str 2668 * @param name 2669 * @return string 2670 */ 2671 public static String notEmpty(String str, String name) { 2672 if (str == null) { 2673 throw new IllegalArgumentException(name + " cannot be null"); 2674 } 2675 if (str.length() == 0) { 2676 throw new IllegalArgumentException(name + " cannot be empty"); 2677 } 2678 return str; 2679 } 2680 2681 /** 2682 * Check if the object is not null. 2683 * 2684 * @param <T> 2685 * @param obj 2686 * @param name 2687 * @return string 2688 */ 2689 public static <T> T notNull(T obj, String name) { 2690 if (obj == null) { 2691 throw new IllegalArgumentException(name + " cannot be null"); 2692 } 2693 return obj; 2694 } 2695 2696}