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