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