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.cli; 020 021import com.google.common.annotations.VisibleForTesting; 022import org.apache.commons.cli.CommandLine; 023import org.apache.commons.cli.Option; 024import org.apache.commons.cli.OptionBuilder; 025import org.apache.commons.cli.OptionGroup; 026import org.apache.commons.cli.Options; 027import org.apache.commons.cli.ParseException; 028import org.apache.oozie.BuildInfo; 029import org.apache.oozie.client.AuthOozieClient; 030import org.apache.oozie.client.BulkResponse; 031import org.apache.oozie.client.BundleJob; 032import org.apache.oozie.client.CoordinatorAction; 033import org.apache.oozie.client.CoordinatorJob; 034import org.apache.oozie.client.OozieClient; 035import org.apache.oozie.client.OozieClient.SYSTEM_MODE; 036import org.apache.oozie.client.OozieClientException; 037import org.apache.oozie.client.WorkflowAction; 038import org.apache.oozie.client.WorkflowJob; 039import org.apache.oozie.client.XOozieClient; 040import org.apache.oozie.client.rest.JsonTags; 041import org.apache.oozie.client.rest.JsonToBean; 042import org.apache.oozie.client.rest.RestConstants; 043import org.json.simple.JSONArray; 044import org.json.simple.JSONObject; 045import org.w3c.dom.DOMException; 046import org.w3c.dom.Document; 047import org.w3c.dom.Element; 048import org.w3c.dom.Node; 049import org.w3c.dom.NodeList; 050import org.w3c.dom.Text; 051import org.xml.sax.SAXException; 052 053import javax.xml.XMLConstants; 054import javax.xml.parsers.DocumentBuilder; 055import javax.xml.parsers.DocumentBuilderFactory; 056import javax.xml.parsers.ParserConfigurationException; 057import javax.xml.transform.stream.StreamSource; 058import javax.xml.validation.Schema; 059import javax.xml.validation.SchemaFactory; 060import javax.xml.validation.Validator; 061import java.io.File; 062import java.io.FileInputStream; 063import java.io.FileReader; 064import java.io.IOException; 065import java.io.InputStream; 066import java.io.PrintStream; 067import java.text.SimpleDateFormat; 068import java.util.ArrayList; 069import java.util.Date; 070import java.util.List; 071import java.util.Locale; 072import java.util.Map; 073import java.util.Properties; 074import java.util.TimeZone; 075import java.util.TreeMap; 076import java.util.concurrent.Callable; 077import java.util.regex.Matcher; 078import java.util.regex.Pattern; 079 080/** 081 * Oozie command line utility. 082 */ 083public class OozieCLI { 084 public static final String ENV_OOZIE_URL = "OOZIE_URL"; 085 public static final String ENV_OOZIE_DEBUG = "OOZIE_DEBUG"; 086 public static final String ENV_OOZIE_TIME_ZONE = "OOZIE_TIMEZONE"; 087 public static final String ENV_OOZIE_AUTH = "OOZIE_AUTH"; 088 public static final String OOZIE_RETRY_COUNT = "oozie.connection.retry.count"; 089 public static final String WS_HEADER_PREFIX = "header:"; 090 091 public static final String HELP_CMD = "help"; 092 public static final String VERSION_CMD = "version"; 093 public static final String JOB_CMD = "job"; 094 public static final String JOBS_CMD = "jobs"; 095 public static final String ADMIN_CMD = "admin"; 096 public static final String VALIDATE_CMD = "validate"; 097 public static final String SLA_CMD = "sla"; 098 public static final String PIG_CMD = "pig"; 099 public static final String HIVE_CMD = "hive"; 100 public static final String SQOOP_CMD = "sqoop"; 101 public static final String MR_CMD = "mapreduce"; 102 public static final String INFO_CMD = "info"; 103 104 public static final String OOZIE_OPTION = "oozie"; 105 public static final String CONFIG_OPTION = "config"; 106 public static final String SUBMIT_OPTION = "submit"; 107 public static final String OFFSET_OPTION = "offset"; 108 public static final String START_OPTION = "start"; 109 public static final String RUN_OPTION = "run"; 110 public static final String DRYRUN_OPTION = "dryrun"; 111 public static final String SUSPEND_OPTION = "suspend"; 112 public static final String RESUME_OPTION = "resume"; 113 public static final String KILL_OPTION = "kill"; 114 public static final String CHANGE_OPTION = "change"; 115 public static final String CHANGE_VALUE_OPTION = "value"; 116 public static final String RERUN_OPTION = "rerun"; 117 public static final String INFO_OPTION = "info"; 118 public static final String LOG_OPTION = "log"; 119 public static final String ACTION_OPTION = "action"; 120 public static final String DEFINITION_OPTION = "definition"; 121 public static final String CONFIG_CONTENT_OPTION = "configcontent"; 122 public static final String SQOOP_COMMAND_OPTION = "command"; 123 public static final String SHOWDIFF_OPTION = "diff"; 124 public static final String UPDATE_OPTION = "update"; 125 public static final String IGNORE_OPTION = "ignore"; 126 public static final String POLL_OPTION = "poll"; 127 public static final String TIMEOUT_OPTION = "timeout"; 128 public static final String INTERVAL_OPTION = "interval"; 129 130 public static final String DO_AS_OPTION = "doas"; 131 132 public static final String LEN_OPTION = "len"; 133 public static final String FILTER_OPTION = "filter"; 134 public static final String JOBTYPE_OPTION = "jobtype"; 135 public static final String SYSTEM_MODE_OPTION = "systemmode"; 136 public static final String VERSION_OPTION = "version"; 137 public static final String STATUS_OPTION = "status"; 138 public static final String LOCAL_TIME_OPTION = "localtime"; 139 public static final String TIME_ZONE_OPTION = "timezone"; 140 public static final String QUEUE_DUMP_OPTION = "queuedump"; 141 public static final String RERUN_COORD_OPTION = "coordinator"; 142 public static final String DATE_OPTION = "date"; 143 public static final String RERUN_REFRESH_OPTION = "refresh"; 144 public static final String RERUN_NOCLEANUP_OPTION = "nocleanup"; 145 public static final String RERUN_FAILED_OPTION = "failed"; 146 public static final String ORDER_OPTION = "order"; 147 148 public static final String UPDATE_SHARELIB_OPTION = "sharelibupdate"; 149 150 public static final String LIST_SHARELIB_LIB_OPTION = "shareliblist"; 151 152 public static final String SERVER_CONFIGURATION_OPTION = "configuration"; 153 public static final String SERVER_OS_ENV_OPTION = "osenv"; 154 public static final String SERVER_JAVA_SYSTEM_PROPERTIES_OPTION = "javasysprops"; 155 156 public static final String METRICS_OPTION = "metrics"; 157 public static final String INSTRUMENTATION_OPTION = "instrumentation"; 158 159 public static final String AUTH_OPTION = "auth"; 160 161 public static final String VERBOSE_OPTION = "verbose"; 162 public static final String VERBOSE_DELIMITER = "\t"; 163 public static final String DEBUG_OPTION = "debug"; 164 165 public static final String SCRIPTFILE_OPTION = "file"; 166 167 public static final String INFO_TIME_ZONES_OPTION = "timezones"; 168 169 public static final String BULK_OPTION = "bulk"; 170 171 public static final String AVAILABLE_SERVERS_OPTION = "servers"; 172 173 public static final String ALL_WORKFLOWS_FOR_COORD_ACTION = "allruns"; 174 175 private static final String[] OOZIE_HELP = { 176 "the env variable '" + ENV_OOZIE_URL + "' is used as default value for the '-" + OOZIE_OPTION + "' option", 177 "the env variable '" + ENV_OOZIE_TIME_ZONE + "' is used as default value for the '-" + TIME_ZONE_OPTION + "' option", 178 "the env variable '" + ENV_OOZIE_AUTH + "' is used as default value for the '-" + AUTH_OPTION + "' option", 179 "custom headers for Oozie web services can be specified using '-D" + WS_HEADER_PREFIX + "NAME=VALUE'" }; 180 181 private static final String RULER; 182 private static final int LINE_WIDTH = 132; 183 184 private static final int RETRY_COUNT = 4; 185 186 private boolean used; 187 188 private static final String INSTANCE_SEPARATOR = "#"; 189 190 private static final String MAPRED_MAPPER = "mapred.mapper.class"; 191 private static final String MAPRED_MAPPER_2 = "mapreduce.map.class"; 192 private static final String MAPRED_REDUCER = "mapred.reducer.class"; 193 private static final String MAPRED_REDUCER_2 = "mapreduce.reduce.class"; 194 private static final String MAPRED_INPUT = "mapred.input.dir"; 195 private static final String MAPRED_OUTPUT = "mapred.output.dir"; 196 197 private static final Pattern GMT_OFFSET_SHORTEN_PATTERN = Pattern.compile("(.* )GMT((?:-|\\+)\\d{2}:\\d{2})"); 198 199 static { 200 StringBuilder sb = new StringBuilder(); 201 for (int i = 0; i < LINE_WIDTH; i++) { 202 sb.append("-"); 203 } 204 RULER = sb.toString(); 205 } 206 207 /** 208 * Entry point for the Oozie CLI when invoked from the command line. 209 * <p/> 210 * Upon completion this method exits the JVM with '0' (success) or '-1' (failure). 211 * 212 * @param args options and arguments for the Oozie CLI. 213 */ 214 public static void main(String[] args) { 215 if (!System.getProperties().containsKey(AuthOozieClient.USE_AUTH_TOKEN_CACHE_SYS_PROP)) { 216 System.setProperty(AuthOozieClient.USE_AUTH_TOKEN_CACHE_SYS_PROP, "true"); 217 } 218 System.exit(new OozieCLI().run(args)); 219 } 220 221 /** 222 * Create an Oozie CLI instance. 223 */ 224 public OozieCLI() { 225 used = false; 226 } 227 228 /** 229 * Return Oozie CLI top help lines. 230 * 231 * @return help lines. 232 */ 233 protected String[] getCLIHelp() { 234 return OOZIE_HELP; 235 } 236 237 /** 238 * Add authentication specific options to oozie cli 239 * 240 * @param options the collection of options to add auth options 241 */ 242 protected void addAuthOptions(Options options) { 243 Option auth = new Option(AUTH_OPTION, true, "select authentication type [SIMPLE|KERBEROS]"); 244 options.addOption(auth); 245 } 246 247 /** 248 * Create option for command line option 'admin' 249 * @return admin options 250 */ 251 protected Options createAdminOptions() { 252 Option oozie = new Option(OOZIE_OPTION, true, "Oozie URL"); 253 Option system_mode = new Option(SYSTEM_MODE_OPTION, true, 254 "Supported in Oozie-2.0 or later versions ONLY. Change oozie system mode [NORMAL|NOWEBSERVICE|SAFEMODE]"); 255 Option status = new Option(STATUS_OPTION, false, "show the current system status"); 256 Option version = new Option(VERSION_OPTION, false, "show Oozie server build version"); 257 Option queuedump = new Option(QUEUE_DUMP_OPTION, false, "show Oozie server queue elements"); 258 Option doAs = new Option(DO_AS_OPTION, true, "doAs user, impersonates as the specified user"); 259 Option availServers = new Option(AVAILABLE_SERVERS_OPTION, false, "list available Oozie servers" 260 + " (more than one only if HA is enabled)"); 261 Option sharelibUpdate = new Option(UPDATE_SHARELIB_OPTION, false, "Update server to use a newer version of sharelib"); 262 Option serverConfiguration = new Option(SERVER_CONFIGURATION_OPTION, false, "show Oozie system configuration"); 263 Option osEnv = new Option(SERVER_OS_ENV_OPTION, false, "show Oozie system OS environment"); 264 Option javaSysProps = new Option(SERVER_JAVA_SYSTEM_PROPERTIES_OPTION, false, "show Oozie Java system properties"); 265 Option metrics = new Option(METRICS_OPTION, false, "show Oozie system metrics"); 266 Option instrumentation = new Option(INSTRUMENTATION_OPTION, false, "show Oozie system instrumentation"); 267 268 Option sharelib = new Option(LIST_SHARELIB_LIB_OPTION, false, 269 "List available sharelib that can be specified in a workflow action"); 270 sharelib.setOptionalArg(true); 271 272 Options adminOptions = new Options(); 273 adminOptions.addOption(oozie); 274 adminOptions.addOption(doAs); 275 OptionGroup group = new OptionGroup(); 276 group.addOption(system_mode); 277 group.addOption(status); 278 group.addOption(version); 279 group.addOption(queuedump); 280 group.addOption(availServers); 281 group.addOption(sharelibUpdate); 282 group.addOption(sharelib); 283 group.addOption(serverConfiguration); 284 group.addOption(osEnv); 285 group.addOption(javaSysProps); 286 group.addOption(metrics); 287 group.addOption(instrumentation); 288 adminOptions.addOptionGroup(group); 289 addAuthOptions(adminOptions); 290 return adminOptions; 291 } 292 293 /** 294 * Create option for command line option 'job' 295 * @return job options 296 */ 297 protected Options createJobOptions() { 298 Option oozie = new Option(OOZIE_OPTION, true, "Oozie URL"); 299 Option config = new Option(CONFIG_OPTION, true, "job configuration file '.xml' or '.properties'"); 300 Option submit = new Option(SUBMIT_OPTION, false, "submit a job"); 301 Option run = new Option(RUN_OPTION, false, "run a job"); 302 Option debug = new Option(DEBUG_OPTION, false, "Use debug mode to see debugging statements on stdout"); 303 Option rerun = new Option(RERUN_OPTION, true, 304 "rerun a job (coordinator requires -action or -date, bundle requires -coordinator or -date)"); 305 Option dryrun = new Option(DRYRUN_OPTION, false, "Dryrun a workflow (since 3.3.2) or coordinator (since 2.0) job without" 306 + " actually executing it"); 307 Option update = new Option(UPDATE_OPTION, true, "Update coord definition and properties"); 308 Option showdiff = new Option(SHOWDIFF_OPTION, true, 309 "Show diff of the new coord definition and properties with the existing one (default true)"); 310 Option start = new Option(START_OPTION, true, "start a job"); 311 Option suspend = new Option(SUSPEND_OPTION, true, "suspend a job"); 312 Option resume = new Option(RESUME_OPTION, true, "resume a job"); 313 Option kill = new Option(KILL_OPTION, true, "kill a job (coordinator can mention -action or -date)"); 314 Option change = new Option(CHANGE_OPTION, true, "change a coordinator or bundle job"); 315 Option changeValue = new Option(CHANGE_VALUE_OPTION, true, 316 "new endtime/concurrency/pausetime value for changing a coordinator job"); 317 Option info = new Option(INFO_OPTION, true, "info of a job"); 318 Option poll = new Option(POLL_OPTION, true, "poll Oozie until a job reaches a terminal state or a timeout occurs"); 319 Option offset = new Option(OFFSET_OPTION, true, "job info offset of actions (default '1', requires -info)"); 320 Option len = new Option(LEN_OPTION, true, "number of actions (default TOTAL ACTIONS, requires -info)"); 321 Option filter = new Option(FILTER_OPTION, true, 322 "<key><comparator><value>[;<key><comparator><value>]*\n" 323 + "(All Coordinator actions satisfying the filters will be retreived).\n" 324 + "key: status or nominaltime\n" 325 + "comparator: =, !=, <, <=, >, >=. = is used as OR and others as AND\n" 326 + "status: values are valid status like SUCCEEDED, KILLED etc. Only = and != apply for status\n" 327 + "nominaltime: time of format yyyy-MM-dd'T'HH:mm'Z'"); 328 Option order = new Option(ORDER_OPTION, true, 329 "order to show coord actions (default ascending order, 'desc' for descending order, requires -info)"); 330 Option localtime = new Option(LOCAL_TIME_OPTION, false, "use local time (same as passing your time zone to -" + 331 TIME_ZONE_OPTION + "). Overrides -" + TIME_ZONE_OPTION + " option"); 332 Option timezone = new Option(TIME_ZONE_OPTION, true, 333 "use time zone with the specified ID (default GMT).\nSee 'oozie info -timezones' for a list"); 334 Option log = new Option(LOG_OPTION, true, "job log"); 335 Option logFilter = new Option( 336 RestConstants.LOG_FILTER_OPTION, true, 337 "job log search parameter. Can be specified as -logfilter opt1=val1;opt2=val1;opt3=val1. " 338 + "Supported options are recent, start, end, loglevel, text, limit and debug"); 339 Option definition = new Option(DEFINITION_OPTION, true, "job definition"); 340 Option config_content = new Option(CONFIG_CONTENT_OPTION, true, "job configuration"); 341 Option verbose = new Option(VERBOSE_OPTION, false, "verbose mode"); 342 Option action = new Option(ACTION_OPTION, true, 343 "coordinator rerun/kill on action ids (requires -rerun/-kill); coordinator log retrieval on action ids" 344 + "(requires -log)"); 345 Option date = new Option(DATE_OPTION, true, 346 "coordinator/bundle rerun on action dates (requires -rerun); coordinator log retrieval on action dates (requires -log)"); 347 Option rerun_coord = new Option(RERUN_COORD_OPTION, true, "bundle rerun on coordinator names (requires -rerun)"); 348 Option rerun_refresh = new Option(RERUN_REFRESH_OPTION, false, 349 "re-materialize the coordinator rerun actions (requires -rerun)"); 350 Option rerun_nocleanup = new Option(RERUN_NOCLEANUP_OPTION, false, 351 "do not clean up output-events of the coordiantor rerun actions (requires -rerun)"); 352 Option rerun_failed = new Option(RERUN_FAILED_OPTION, false, 353 "runs the failed workflow actions of the coordinator actions (requires -rerun)"); 354 Option property = OptionBuilder.withArgName("property=value").hasArgs(2).withValueSeparator().withDescription( 355 "set/override value for given property").create("D"); 356 Option getAllWorkflows = new Option(ALL_WORKFLOWS_FOR_COORD_ACTION, false, 357 "Get workflow jobs corresponding to a coordinator action including all the reruns"); 358 Option ignore = new Option(IGNORE_OPTION, true, 359 "change status of a coordinator job or action to IGNORED" 360 + " (-action required to ignore coord actions)"); 361 Option timeout = new Option(TIMEOUT_OPTION, true, "timeout in minutes (default is 30, negative values indicate no " 362 + "timeout, requires -poll)"); 363 timeout.setType(Integer.class); 364 Option interval = new Option(INTERVAL_OPTION, true, "polling interval in minutes (default is 5, requires -poll)"); 365 interval.setType(Integer.class); 366 367 Option doAs = new Option(DO_AS_OPTION, true, "doAs user, impersonates as the specified user"); 368 369 OptionGroup actions = new OptionGroup(); 370 actions.addOption(submit); 371 actions.addOption(start); 372 actions.addOption(run); 373 actions.addOption(dryrun); 374 actions.addOption(suspend); 375 actions.addOption(resume); 376 actions.addOption(kill); 377 actions.addOption(change); 378 actions.addOption(update); 379 actions.addOption(info); 380 actions.addOption(rerun); 381 actions.addOption(log); 382 actions.addOption(definition); 383 actions.addOption(config_content); 384 actions.addOption(ignore); 385 actions.addOption(poll); 386 actions.setRequired(true); 387 Options jobOptions = new Options(); 388 jobOptions.addOption(oozie); 389 jobOptions.addOption(doAs); 390 jobOptions.addOption(config); 391 jobOptions.addOption(property); 392 jobOptions.addOption(changeValue); 393 jobOptions.addOption(localtime); 394 jobOptions.addOption(timezone); 395 jobOptions.addOption(verbose); 396 jobOptions.addOption(debug); 397 jobOptions.addOption(offset); 398 jobOptions.addOption(len); 399 jobOptions.addOption(filter); 400 jobOptions.addOption(order); 401 jobOptions.addOption(action); 402 jobOptions.addOption(date); 403 jobOptions.addOption(rerun_coord); 404 jobOptions.addOption(rerun_refresh); 405 jobOptions.addOption(rerun_nocleanup); 406 jobOptions.addOption(rerun_failed); 407 jobOptions.addOption(getAllWorkflows); 408 jobOptions.addOptionGroup(actions); 409 jobOptions.addOption(logFilter); 410 jobOptions.addOption(timeout); 411 jobOptions.addOption(interval); 412 addAuthOptions(jobOptions); 413 jobOptions.addOption(showdiff); 414 415 //Needed to make dryrun and update mutually exclusive options 416 OptionGroup updateOption = new OptionGroup(); 417 updateOption.addOption(dryrun); 418 jobOptions.addOptionGroup(updateOption); 419 return jobOptions; 420 } 421 422 /** 423 * Create option for command line option 'jobs' 424 * @return jobs options 425 */ 426 protected Options createJobsOptions() { 427 Option oozie = new Option(OOZIE_OPTION, true, "Oozie URL"); 428 Option start = new Option(OFFSET_OPTION, true, "jobs offset (default '1')"); 429 Option jobtype = new Option(JOBTYPE_OPTION, true, 430 "job type ('Supported in Oozie-2.0 or later versions ONLY - 'coordinator' or 'bundle' or 'wf'(default))"); 431 Option len = new Option(LEN_OPTION, true, "number of jobs (default '100')"); 432 Option filter = new Option(FILTER_OPTION, true, 433 "user=<U>\\;name=<N>\\;group=<G>\\;status=<S>\\;frequency=<F>\\;unit=<M>" + 434 "\\;startcreatedtime=<SC>\\;endcreatedtime=<EC> " + 435 "(valid unit values are 'months', 'days', 'hours' or 'minutes'. " + 436 "startcreatedtime, endcreatedtime: time of format yyyy-MM-dd'T'HH:mm'Z')"); 437 Option localtime = new Option(LOCAL_TIME_OPTION, false, "use local time (same as passing your time zone to -" + 438 TIME_ZONE_OPTION + "). Overrides -" + TIME_ZONE_OPTION + " option"); 439 Option kill = new Option(KILL_OPTION, false, "bulk kill operation"); 440 Option suspend = new Option(SUSPEND_OPTION, false, "bulk suspend operation"); 441 Option resume = new Option(RESUME_OPTION, false, "bulk resume operation"); 442 Option timezone = new Option(TIME_ZONE_OPTION, true, 443 "use time zone with the specified ID (default GMT).\nSee 'oozie info -timezones' for a list"); 444 Option verbose = new Option(VERBOSE_OPTION, false, "verbose mode"); 445 Option doAs = new Option(DO_AS_OPTION, true, "doAs user, impersonates as the specified user"); 446 Option bulkMonitor = new Option(BULK_OPTION, true, "key-value pairs to filter bulk jobs response. e.g. bundle=<B>\\;" + 447 "coordinators=<C>\\;actionstatus=<S>\\;startcreatedtime=<SC>\\;endcreatedtime=<EC>\\;" + 448 "startscheduledtime=<SS>\\;endscheduledtime=<ES>\\; bundle, coordinators and actionstatus can be multiple comma separated values" + 449 "bundle and coordinators can be id(s) or appName(s) of those jobs. Specifying bundle is mandatory, other params are optional"); 450 start.setType(Integer.class); 451 len.setType(Integer.class); 452 Options jobsOptions = new Options(); 453 jobsOptions.addOption(oozie); 454 jobsOptions.addOption(doAs); 455 jobsOptions.addOption(localtime); 456 jobsOptions.addOption(kill); 457 jobsOptions.addOption(suspend); 458 jobsOptions.addOption(resume); 459 jobsOptions.addOption(timezone); 460 jobsOptions.addOption(start); 461 jobsOptions.addOption(len); 462 jobsOptions.addOption(oozie); 463 jobsOptions.addOption(filter); 464 jobsOptions.addOption(jobtype); 465 jobsOptions.addOption(verbose); 466 jobsOptions.addOption(bulkMonitor); 467 addAuthOptions(jobsOptions); 468 return jobsOptions; 469 } 470 471 /** 472 * Create option for command line option 'sla' 473 * 474 * @return sla options 475 */ 476 protected Options createSlaOptions() { 477 Option oozie = new Option(OOZIE_OPTION, true, "Oozie URL"); 478 Option start = new Option(OFFSET_OPTION, true, "start offset (default '0')"); 479 Option len = new Option(LEN_OPTION, true, "number of results (default '100', max '1000')"); 480 Option filter = new Option(FILTER_OPTION, true, "filter of SLA events. e.g., jobid=<J>\\;appname=<A>"); 481 start.setType(Integer.class); 482 len.setType(Integer.class); 483 Options slaOptions = new Options(); 484 slaOptions.addOption(start); 485 slaOptions.addOption(len); 486 slaOptions.addOption(filter); 487 slaOptions.addOption(oozie); 488 addAuthOptions(slaOptions); 489 return slaOptions; 490 } 491 492 /** 493 * Create option for command line option 'validate' 494 * 495 * @return validate options 496 */ 497 protected Options createValidateOptions() { 498 Option oozie = new Option(OOZIE_OPTION, true, "Oozie URL"); 499 Options validateOption = new Options(); 500 validateOption.addOption(oozie); 501 addAuthOptions(validateOption); 502 return validateOption; 503 } 504 505 /** 506 * Create option for command line option 'pig' or 'hive' 507 * @return pig or hive options 508 */ 509 @SuppressWarnings("static-access") 510 protected Options createScriptLanguageOptions(String jobType) { 511 Option oozie = new Option(OOZIE_OPTION, true, "Oozie URL"); 512 Option config = new Option(CONFIG_OPTION, true, "job configuration file '.properties'"); 513 Option file = new Option(SCRIPTFILE_OPTION, true, jobType + " script"); 514 Option property = OptionBuilder.withArgName("property=value").hasArgs(2).withValueSeparator().withDescription( 515 "set/override value for given property").create("D"); 516 Option params = OptionBuilder.withArgName("property=value").hasArgs(2).withValueSeparator().withDescription( 517 "set parameters for script").create("P"); 518 Option doAs = new Option(DO_AS_OPTION, true, "doAs user, impersonates as the specified user"); 519 Options Options = new Options(); 520 Options.addOption(oozie); 521 Options.addOption(doAs); 522 Options.addOption(config); 523 Options.addOption(property); 524 Options.addOption(params); 525 Options.addOption(file); 526 addAuthOptions(Options); 527 return Options; 528 } 529 530 /** 531 * Create option for command line option 'sqoop' 532 * @return sqoop options 533 */ 534 @SuppressWarnings("static-access") 535 protected Options createSqoopCLIOptions() { 536 Option oozie = new Option(OOZIE_OPTION, true, "Oozie URL"); 537 Option config = new Option(CONFIG_OPTION, true, "job configuration file '.properties'"); 538 Option command = OptionBuilder.withArgName(SQOOP_COMMAND_OPTION).hasArgs().withValueSeparator().withDescription( 539 "sqoop command").create(SQOOP_COMMAND_OPTION); 540 Option property = OptionBuilder.withArgName("property=value").hasArgs(2).withValueSeparator().withDescription( 541 "set/override value for given property").create("D"); 542 Option doAs = new Option(DO_AS_OPTION, true, "doAs user, impersonates as the specified user"); 543 Options Options = new Options(); 544 Options.addOption(oozie); 545 Options.addOption(doAs); 546 Options.addOption(config); 547 Options.addOption(property); 548 Options.addOption(command); 549 addAuthOptions(Options); 550 return Options; 551 } 552 553 /** 554 * Create option for command line option 'info' 555 * @return info options 556 */ 557 protected Options createInfoOptions() { 558 Option timezones = new Option(INFO_TIME_ZONES_OPTION, false, "display a list of available time zones"); 559 Options infoOptions = new Options(); 560 infoOptions.addOption(timezones); 561 return infoOptions; 562 } 563 564 /** 565 * Create option for command line option 'mapreduce' 566 * @return mapreduce options 567 */ 568 @SuppressWarnings("static-access") 569 protected Options createMROptions() { 570 Option oozie = new Option(OOZIE_OPTION, true, "Oozie URL"); 571 Option config = new Option(CONFIG_OPTION, true, "job configuration file '.properties'"); 572 Option property = OptionBuilder.withArgName("property=value").hasArgs(2).withValueSeparator().withDescription( 573 "set/override value for given property").create("D"); 574 Option doAs = new Option(DO_AS_OPTION, true, "doAs user, impersonates as the specified user"); 575 Options mrOptions = new Options(); 576 mrOptions.addOption(oozie); 577 mrOptions.addOption(doAs); 578 mrOptions.addOption(config); 579 mrOptions.addOption(property); 580 addAuthOptions(mrOptions); 581 return mrOptions; 582 } 583 584 /** 585 * Run a CLI programmatically. 586 * <p/> 587 * It does not exit the JVM. 588 * <p/> 589 * A CLI instance can be used only once. 590 * 591 * @param args options and arguments for the Oozie CLI. 592 * @return '0' (success), '-1' (failure). 593 */ 594 public synchronized int run(String[] args) { 595 if (used) { 596 throw new IllegalStateException("CLI instance already used"); 597 } 598 used = true; 599 final CLIParser parser = getCLIParser(); 600 try { 601 final CLIParser.Command command = parser.parse(args); 602 603 String doAsUser = command.getCommandLine().getOptionValue(DO_AS_OPTION); 604 605 if (doAsUser != null) { 606 OozieClient.doAs(doAsUser, new Callable<Void>() { 607 @Override 608 public Void call() throws Exception { 609 processCommand(parser, command); 610 return null; 611 } 612 }); 613 } 614 else { 615 processCommand(parser, command); 616 } 617 return 0; 618 } 619 catch (OozieCLIException ex) { 620 System.err.println("Error: " + ex.getMessage()); 621 return -1; 622 } 623 catch (ParseException ex) { 624 System.err.println("Invalid sub-command: " + ex.getMessage()); 625 System.err.println(); 626 System.err.println(parser.shortHelp()); 627 return -1; 628 } 629 catch (Exception ex) { 630 ex.printStackTrace(); 631 System.err.println(ex.getMessage()); 632 return -1; 633 } 634 } 635 636 @VisibleForTesting 637 public CLIParser getCLIParser(){ 638 CLIParser parser = new CLIParser(OOZIE_OPTION, getCLIHelp()); 639 parser.addCommand(HELP_CMD, "", "display usage for all commands or specified command", new Options(), false); 640 parser.addCommand(VERSION_CMD, "", "show client version", new Options(), false); 641 parser.addCommand(JOB_CMD, "", "job operations", createJobOptions(), false); 642 parser.addCommand(JOBS_CMD, "", "jobs status", createJobsOptions(), false); 643 parser.addCommand(ADMIN_CMD, "", "admin operations", createAdminOptions(), false); 644 parser.addCommand(VALIDATE_CMD, "", "validate a workflow, coordinator, bundle XML file", createValidateOptions(), true); 645 parser.addCommand(SLA_CMD, "", "sla operations (Deprecated with Oozie 4.0)", createSlaOptions(), false); 646 parser.addCommand(PIG_CMD, "-X ", "submit a pig job, everything after '-X' are pass-through parameters to pig, any '-D' " 647 + "arguments after '-X' are put in <configuration>", createScriptLanguageOptions(PIG_CMD), true); 648 parser.addCommand(HIVE_CMD, "-X ", "submit a hive job, everything after '-X' are pass-through parameters to hive, any '-D' " 649 + "arguments after '-X' are put in <configuration>", createScriptLanguageOptions(HIVE_CMD), true); 650 parser.addCommand(SQOOP_CMD, "-X ", "submit a sqoop job, everything after '-X' are pass-through parameters " + 651 "to sqoop, any '-D' arguments after '-X' are put in <configuration>", createSqoopCLIOptions(), true); 652 parser.addCommand(INFO_CMD, "", "get more detailed info about specific topics", createInfoOptions(), false); 653 parser.addCommand(MR_CMD, "", "submit a mapreduce job", createMROptions(), false); 654 return parser; 655 } 656 657 public void processCommand(CLIParser parser, CLIParser.Command command) throws Exception { 658 if (command.getName().equals(HELP_CMD)) { 659 parser.showHelp(command.getCommandLine()); 660 } 661 else if (command.getName().equals(JOB_CMD)) { 662 jobCommand(command.getCommandLine()); 663 } 664 else if (command.getName().equals(JOBS_CMD)) { 665 jobsCommand(command.getCommandLine()); 666 } 667 else if (command.getName().equals(ADMIN_CMD)) { 668 adminCommand(command.getCommandLine()); 669 } 670 else if (command.getName().equals(VERSION_CMD)) { 671 versionCommand(); 672 } 673 else if (command.getName().equals(VALIDATE_CMD)) { 674 validateCommand(command.getCommandLine()); 675 } 676 else if (command.getName().equals(SLA_CMD)) { 677 slaCommand(command.getCommandLine()); 678 } 679 else if (command.getName().equals(PIG_CMD)) { 680 scriptLanguageCommand(command.getCommandLine(), PIG_CMD); 681 } 682 else if (command.getName().equals(HIVE_CMD)) { 683 scriptLanguageCommand(command.getCommandLine(), HIVE_CMD); 684 } 685 else if (command.getName().equals(SQOOP_CMD)) { 686 sqoopCommand(command.getCommandLine()); 687 } 688 else if (command.getName().equals(INFO_CMD)) { 689 infoCommand(command.getCommandLine()); 690 } 691 else if (command.getName().equals(MR_CMD)){ 692 mrCommand(command.getCommandLine()); 693 } 694 } 695 protected String getOozieUrl(CommandLine commandLine) { 696 String url = commandLine.getOptionValue(OOZIE_OPTION); 697 if (url == null) { 698 url = System.getenv(ENV_OOZIE_URL); 699 if (url == null) { 700 throw new IllegalArgumentException( 701 "Oozie URL is not available neither in command option or in the environment"); 702 } 703 } 704 return url; 705 } 706 707 private String getTimeZoneId(CommandLine commandLine) 708 { 709 if (commandLine.hasOption(LOCAL_TIME_OPTION)) { 710 return null; 711 } 712 if (commandLine.hasOption(TIME_ZONE_OPTION)) { 713 return commandLine.getOptionValue(TIME_ZONE_OPTION); 714 } 715 String timeZoneId = System.getenv(ENV_OOZIE_TIME_ZONE); 716 if (timeZoneId != null) { 717 return timeZoneId; 718 } 719 return "GMT"; 720 } 721 722 // Canibalized from Hadoop <code>Configuration.loadResource()</code>. 723 private Properties parse(InputStream is, Properties conf) throws IOException { 724 try { 725 DocumentBuilderFactory docBuilderFactory = DocumentBuilderFactory.newInstance(); 726 docBuilderFactory.setNamespaceAware(true); 727 // support for includes in the xml file 728 docBuilderFactory.setXIncludeAware(true); 729 // ignore all comments inside the xml file 730 docBuilderFactory.setIgnoringComments(true); 731 docBuilderFactory.setExpandEntityReferences(false); 732 docBuilderFactory.setFeature("http://apache.org/xml/features/disallow-doctype-decl", true); 733 DocumentBuilder builder = docBuilderFactory.newDocumentBuilder(); 734 Document doc = builder.parse(is); 735 return parseDocument(doc, conf); 736 } 737 catch (SAXException e) { 738 throw new IOException(e); 739 } 740 catch (ParserConfigurationException e) { 741 throw new IOException(e); 742 } 743 } 744 745 // Canibalized from Hadoop <code>Configuration.loadResource()</code>. 746 private Properties parseDocument(Document doc, Properties conf) throws IOException { 747 try { 748 Element root = doc.getDocumentElement(); 749 if (!"configuration".equals(root.getLocalName())) { 750 throw new RuntimeException("bad conf file: top-level element not <configuration>"); 751 } 752 NodeList props = root.getChildNodes(); 753 for (int i = 0; i < props.getLength(); i++) { 754 Node propNode = props.item(i); 755 if (!(propNode instanceof Element)) { 756 continue; 757 } 758 Element prop = (Element) propNode; 759 if (!"property".equals(prop.getLocalName())) { 760 throw new RuntimeException("bad conf file: element not <property>"); 761 } 762 NodeList fields = prop.getChildNodes(); 763 String attr = null; 764 String value = null; 765 for (int j = 0; j < fields.getLength(); j++) { 766 Node fieldNode = fields.item(j); 767 if (!(fieldNode instanceof Element)) { 768 continue; 769 } 770 Element field = (Element) fieldNode; 771 if ("name".equals(field.getLocalName()) && field.hasChildNodes()) { 772 attr = ((Text) field.getFirstChild()).getData(); 773 } 774 if ("value".equals(field.getLocalName()) && field.hasChildNodes()) { 775 value = ((Text) field.getFirstChild()).getData(); 776 } 777 } 778 779 if (attr != null && value != null) { 780 conf.setProperty(attr, value); 781 } 782 } 783 return conf; 784 } 785 catch (DOMException e) { 786 throw new IOException(e); 787 } 788 } 789 790 private Properties getConfiguration(OozieClient wc, CommandLine commandLine) throws IOException { 791 if (!isConfigurationSpecified(wc, commandLine)) { 792 throw new IOException("configuration is not specified"); 793 } 794 Properties conf = wc.createConfiguration(); 795 String configFile = commandLine.getOptionValue(CONFIG_OPTION); 796 if (configFile != null) { 797 File file = new File(configFile); 798 if (!file.exists()) { 799 throw new IOException("configuration file [" + configFile + "] not found"); 800 } 801 if (configFile.endsWith(".properties")) { 802 conf.load(new FileReader(file)); 803 } 804 else if (configFile.endsWith(".xml")) { 805 parse(new FileInputStream(configFile), conf); 806 } 807 else { 808 throw new IllegalArgumentException("configuration must be a '.properties' or a '.xml' file"); 809 } 810 } 811 if (commandLine.hasOption("D")) { 812 Properties commandLineProperties = commandLine.getOptionProperties("D"); 813 conf.putAll(commandLineProperties); 814 } 815 return conf; 816 } 817 818 /** 819 * Check if configuration has specified 820 * @param wc 821 * @param commandLine 822 * @return 823 * @throws IOException 824 */ 825 private boolean isConfigurationSpecified(OozieClient wc, CommandLine commandLine) throws IOException { 826 boolean isConf = false; 827 String configFile = commandLine.getOptionValue(CONFIG_OPTION); 828 if (configFile == null) { 829 isConf = false; 830 } 831 else { 832 isConf = new File(configFile).exists(); 833 } 834 if (commandLine.hasOption("D")) { 835 isConf = true; 836 } 837 return isConf; 838 } 839 840 /** 841 * @param commandLine command line string. 842 * @return change value specified by -value. 843 * @throws OozieCLIException 844 */ 845 private String getChangeValue(CommandLine commandLine) throws OozieCLIException { 846 String changeValue = commandLine.getOptionValue(CHANGE_VALUE_OPTION); 847 848 if (changeValue == null) { 849 throw new OozieCLIException("-value option needs to be specified for -change option"); 850 } 851 852 return changeValue; 853 } 854 855 protected void addHeader(OozieClient wc) { 856 for (Map.Entry entry : System.getProperties().entrySet()) { 857 String key = (String) entry.getKey(); 858 if (key.startsWith(WS_HEADER_PREFIX)) { 859 String header = key.substring(WS_HEADER_PREFIX.length()); 860 wc.setHeader(header, (String) entry.getValue()); 861 } 862 } 863 } 864 865 /** 866 * Get auth option from command line 867 * 868 * @param commandLine the command line object 869 * @return auth option 870 */ 871 protected String getAuthOption(CommandLine commandLine) { 872 String authOpt = commandLine.getOptionValue(AUTH_OPTION); 873 if (authOpt == null) { 874 authOpt = System.getenv(ENV_OOZIE_AUTH); 875 } 876 if (commandLine.hasOption(DEBUG_OPTION)) { 877 System.out.println(" Auth type : " + authOpt); 878 } 879 return authOpt; 880 } 881 882 /** 883 * Create a OozieClient. 884 * <p/> 885 * It injects any '-Dheader:' as header to the the {@link org.apache.oozie.client.OozieClient}. 886 * 887 * @param commandLine the parsed command line options. 888 * @return a pre configured eXtended workflow client. 889 * @throws OozieCLIException thrown if the OozieClient could not be configured. 890 */ 891 protected OozieClient createOozieClient(CommandLine commandLine) throws OozieCLIException { 892 return createXOozieClient(commandLine); 893 } 894 895 /** 896 * Create a XOozieClient. 897 * <p/> 898 * It injects any '-Dheader:' as header to the the {@link org.apache.oozie.client.OozieClient}. 899 * 900 * @param commandLine the parsed command line options. 901 * @return a pre configured eXtended workflow client. 902 * @throws OozieCLIException thrown if the XOozieClient could not be configured. 903 */ 904 protected XOozieClient createXOozieClient(CommandLine commandLine) throws OozieCLIException { 905 XOozieClient wc = new AuthOozieClient(getOozieUrl(commandLine), getAuthOption(commandLine)); 906 addHeader(wc); 907 setDebugMode(wc,commandLine.hasOption(DEBUG_OPTION)); 908 setRetryCount(wc); 909 return wc; 910 } 911 912 protected void setDebugMode(OozieClient wc, boolean debugOpt) { 913 914 String debug = System.getenv(ENV_OOZIE_DEBUG); 915 if (debug != null && !debug.isEmpty()) { 916 int debugVal = 0; 917 try { 918 debugVal = Integer.parseInt(debug.trim()); 919 } 920 catch (Exception ex) { 921 System.out.println("Unable to parse the debug settings. May be not an integer [" + debug + "]"); 922 ex.printStackTrace(); 923 } 924 wc.setDebugMode(debugVal); 925 } 926 else if(debugOpt){ // CLI argument "-debug" used 927 wc.setDebugMode(1); 928 } 929 } 930 931 protected void setRetryCount(OozieClient wc) { 932 String retryCount = System.getProperty(OOZIE_RETRY_COUNT); 933 if (retryCount != null && !retryCount.isEmpty()) { 934 try { 935 int retry = Integer.parseInt(retryCount.trim()); 936 wc.setRetryCount(retry); 937 } 938 catch (Exception ex) { 939 System.err.println("Unable to parse the retry settings. May be not an integer [" + retryCount + "]"); 940 ex.printStackTrace(); 941 } 942 } 943 } 944 945 private static String JOB_ID_PREFIX = "job: "; 946 947 private void jobCommand(CommandLine commandLine) throws IOException, OozieCLIException { 948 XOozieClient wc = createXOozieClient(commandLine); 949 950 List<String> options = new ArrayList<String>(); 951 for (Option option : commandLine.getOptions()) { 952 options.add(option.getOpt()); 953 } 954 955 try { 956 if (options.contains(SUBMIT_OPTION)) { 957 System.out.println(JOB_ID_PREFIX + wc.submit(getConfiguration(wc, commandLine))); 958 } 959 else if (options.contains(START_OPTION)) { 960 wc.start(commandLine.getOptionValue(START_OPTION)); 961 } 962 else if (options.contains(DRYRUN_OPTION) && !options.contains(UPDATE_OPTION)) { 963 String dryrunStr = wc.dryrun(getConfiguration(wc, commandLine)); 964 if (dryrunStr.equals("OK")) { // workflow 965 System.out.println("OK"); 966 } else { // coordinator 967 String[] dryrunStrs = dryrunStr.split("action for new instance"); 968 int arraysize = dryrunStrs.length; 969 System.out.println("***coordJob after parsing: ***"); 970 System.out.println(dryrunStrs[0]); 971 int aLen = dryrunStrs.length - 1; 972 if (aLen < 0) { 973 aLen = 0; 974 } 975 System.out.println("***total coord actions is " + aLen + " ***"); 976 for (int i = 1; i <= arraysize - 1; i++) { 977 System.out.println(RULER); 978 System.out.println("coordAction instance: " + i + ":"); 979 System.out.println(dryrunStrs[i]); 980 } 981 } 982 } 983 else if (options.contains(SUSPEND_OPTION)) { 984 wc.suspend(commandLine.getOptionValue(SUSPEND_OPTION)); 985 } 986 else if (options.contains(RESUME_OPTION)) { 987 wc.resume(commandLine.getOptionValue(RESUME_OPTION)); 988 } 989 else if (options.contains(IGNORE_OPTION)) { 990 String ignoreScope = null; 991 if (options.contains(ACTION_OPTION)) { 992 ignoreScope = commandLine.getOptionValue(ACTION_OPTION); 993 if (ignoreScope == null || ignoreScope.isEmpty()) { 994 throw new OozieCLIException("-" + ACTION_OPTION + " is empty"); 995 } 996 } 997 printCoordActionsStatus(wc.ignore(commandLine.getOptionValue(IGNORE_OPTION), ignoreScope)); 998 } 999 else if (options.contains(KILL_OPTION)) { 1000 if (commandLine.getOptionValue(KILL_OPTION).contains("-C") 1001 && (options.contains(DATE_OPTION) || options.contains(ACTION_OPTION))) { 1002 String coordJobId = commandLine.getOptionValue(KILL_OPTION); 1003 String scope = null; 1004 String rangeType = null; 1005 if (options.contains(DATE_OPTION) && options.contains(ACTION_OPTION)) { 1006 throw new OozieCLIException("Invalid options provided for rerun: either" + DATE_OPTION + " or " 1007 + ACTION_OPTION + " expected. Don't use both at the same time."); 1008 } 1009 if (options.contains(DATE_OPTION)) { 1010 rangeType = RestConstants.JOB_COORD_SCOPE_DATE; 1011 scope = commandLine.getOptionValue(DATE_OPTION); 1012 } 1013 else if (options.contains(ACTION_OPTION)) { 1014 rangeType = RestConstants.JOB_COORD_SCOPE_ACTION; 1015 scope = commandLine.getOptionValue(ACTION_OPTION); 1016 } 1017 else { 1018 throw new OozieCLIException("Invalid options provided for rerun: " + DATE_OPTION + " or " 1019 + ACTION_OPTION + " expected."); 1020 } 1021 printCoordActions(wc.kill(coordJobId, rangeType, scope)); 1022 } 1023 else { 1024 wc.kill(commandLine.getOptionValue(KILL_OPTION)); 1025 } 1026 } 1027 else if (options.contains(CHANGE_OPTION)) { 1028 wc.change(commandLine.getOptionValue(CHANGE_OPTION), getChangeValue(commandLine)); 1029 } 1030 else if (options.contains(RUN_OPTION)) { 1031 System.out.println(JOB_ID_PREFIX + wc.run(getConfiguration(wc, commandLine))); 1032 } 1033 else if (options.contains(RERUN_OPTION)) { 1034 if (commandLine.getOptionValue(RERUN_OPTION).contains("-W")) { 1035 if (isConfigurationSpecified(wc, commandLine)) { 1036 wc.reRun(commandLine.getOptionValue(RERUN_OPTION), getConfiguration(wc, commandLine)); 1037 } 1038 else { 1039 wc.reRun(commandLine.getOptionValue(RERUN_OPTION), new Properties()); 1040 } 1041 } 1042 else if (commandLine.getOptionValue(RERUN_OPTION).contains("-B")) { 1043 String bundleJobId = commandLine.getOptionValue(RERUN_OPTION); 1044 String coordScope = null; 1045 String dateScope = null; 1046 boolean refresh = false; 1047 boolean noCleanup = false; 1048 if (options.contains(ACTION_OPTION)) { 1049 throw new OozieCLIException("Invalid options provided for bundle rerun. " + ACTION_OPTION 1050 + " is not valid for bundle rerun"); 1051 } 1052 if (options.contains(DATE_OPTION)) { 1053 dateScope = commandLine.getOptionValue(DATE_OPTION); 1054 } 1055 1056 if (options.contains(RERUN_COORD_OPTION)) { 1057 coordScope = commandLine.getOptionValue(RERUN_COORD_OPTION); 1058 } 1059 1060 if (options.contains(RERUN_REFRESH_OPTION)) { 1061 refresh = true; 1062 } 1063 if (options.contains(RERUN_NOCLEANUP_OPTION)) { 1064 noCleanup = true; 1065 } 1066 wc.reRunBundle(bundleJobId, coordScope, dateScope, refresh, noCleanup); 1067 if (coordScope != null && !coordScope.isEmpty()) { 1068 System.out.println("Coordinators [" + coordScope + "] of bundle " + bundleJobId 1069 + " are scheduled to rerun on date ranges [" + dateScope + "]."); 1070 } 1071 else { 1072 System.out.println("All coordinators of bundle " + bundleJobId 1073 + " are scheduled to rerun on the date ranges [" + dateScope + "]."); 1074 } 1075 } 1076 else { 1077 String coordJobId = commandLine.getOptionValue(RERUN_OPTION); 1078 String scope = null; 1079 String rerunType = null; 1080 boolean refresh = false; 1081 boolean noCleanup = false; 1082 if (options.contains(DATE_OPTION) && options.contains(ACTION_OPTION)) { 1083 throw new OozieCLIException("Invalid options provided for rerun: either" + DATE_OPTION + " or " 1084 + ACTION_OPTION + " expected. Don't use both at the same time."); 1085 } 1086 if (options.contains(DATE_OPTION)) { 1087 rerunType = RestConstants.JOB_COORD_SCOPE_DATE; 1088 scope = commandLine.getOptionValue(DATE_OPTION); 1089 } 1090 else if (options.contains(ACTION_OPTION)) { 1091 rerunType = RestConstants.JOB_COORD_SCOPE_ACTION; 1092 scope = commandLine.getOptionValue(ACTION_OPTION); 1093 } 1094 else { 1095 throw new OozieCLIException("Invalid options provided for rerun: " + DATE_OPTION + " or " 1096 + ACTION_OPTION + " expected."); 1097 } 1098 if (options.contains(RERUN_REFRESH_OPTION)) { 1099 refresh = true; 1100 } 1101 if (options.contains(RERUN_NOCLEANUP_OPTION)) { 1102 noCleanup = true; 1103 } 1104 if (options.contains(RERUN_FAILED_OPTION)) { 1105 printCoordActions(wc.reRunCoord(coordJobId, rerunType, scope, refresh, noCleanup, true)); 1106 } else { 1107 printCoordActions(wc.reRunCoord(coordJobId, rerunType, scope, refresh, noCleanup)); 1108 } 1109 } 1110 } 1111 else if (options.contains(INFO_OPTION)) { 1112 String timeZoneId = getTimeZoneId(commandLine); 1113 final String optionValue = commandLine.getOptionValue(INFO_OPTION); 1114 if (optionValue.endsWith("-B")) { 1115 String filter = commandLine.getOptionValue(FILTER_OPTION); 1116 if (filter != null) { 1117 throw new OozieCLIException("Filter option is currently not supported for a Bundle job"); 1118 } 1119 printBundleJob(wc.getBundleJobInfo(optionValue), timeZoneId, 1120 options.contains(VERBOSE_OPTION)); 1121 } 1122 else if (optionValue.endsWith("-C")) { 1123 String s = commandLine.getOptionValue(OFFSET_OPTION); 1124 int start = Integer.parseInt((s != null) ? s : "-1"); 1125 s = commandLine.getOptionValue(LEN_OPTION); 1126 int len = Integer.parseInt((s != null) ? s : "-1"); 1127 String filter = commandLine.getOptionValue(FILTER_OPTION); 1128 String order = commandLine.getOptionValue(ORDER_OPTION); 1129 printCoordJob(wc.getCoordJobInfo(optionValue, filter, start, len, order), timeZoneId, 1130 options.contains(VERBOSE_OPTION)); 1131 } 1132 else if (optionValue.contains("-C@")) { 1133 if (options.contains(ALL_WORKFLOWS_FOR_COORD_ACTION)) { 1134 printWfsForCoordAction(wc.getWfsForCoordAction(optionValue), timeZoneId); 1135 } 1136 else { 1137 String filter = commandLine.getOptionValue(FILTER_OPTION); 1138 if (filter != null) { 1139 throw new OozieCLIException("Filter option is not supported for a Coordinator action"); 1140 } 1141 printCoordAction(wc.getCoordActionInfo(optionValue), timeZoneId); 1142 } 1143 } 1144 else if (optionValue.contains("-W@")) { 1145 String filter = commandLine.getOptionValue(FILTER_OPTION); 1146 if (filter != null) { 1147 throw new OozieCLIException("Filter option is not supported for a Workflow action"); 1148 } 1149 printWorkflowAction(wc.getWorkflowActionInfo(optionValue), timeZoneId, 1150 options.contains(VERBOSE_OPTION)); 1151 } 1152 else { 1153 String filter = commandLine.getOptionValue(FILTER_OPTION); 1154 if (filter != null) { 1155 throw new OozieCLIException("Filter option is currently not supported for a Workflow job"); 1156 } 1157 String s = commandLine.getOptionValue(OFFSET_OPTION); 1158 int start = Integer.parseInt((s != null) ? s : "0"); 1159 s = commandLine.getOptionValue(LEN_OPTION); 1160 String jobtype = commandLine.getOptionValue(JOBTYPE_OPTION); 1161 jobtype = (jobtype != null) ? jobtype : "wf"; 1162 int len = Integer.parseInt((s != null) ? s : "0"); 1163 printJob(wc.getJobInfo(optionValue, start, len), timeZoneId, 1164 options.contains(VERBOSE_OPTION)); 1165 } 1166 } 1167 else if (options.contains(LOG_OPTION)) { 1168 PrintStream ps = System.out; 1169 String logFilter = null; 1170 if (options.contains(RestConstants.LOG_FILTER_OPTION)) { 1171 logFilter = commandLine.getOptionValue(RestConstants.LOG_FILTER_OPTION); 1172 } 1173 if (commandLine.getOptionValue(LOG_OPTION).contains("-C")) { 1174 String logRetrievalScope = null; 1175 String logRetrievalType = null; 1176 if (options.contains(ACTION_OPTION)) { 1177 logRetrievalType = RestConstants.JOB_LOG_ACTION; 1178 logRetrievalScope = commandLine.getOptionValue(ACTION_OPTION); 1179 } 1180 if (options.contains(DATE_OPTION)) { 1181 logRetrievalType = RestConstants.JOB_LOG_DATE; 1182 logRetrievalScope = commandLine.getOptionValue(DATE_OPTION); 1183 } 1184 try { 1185 wc.getJobLog(commandLine.getOptionValue(LOG_OPTION), logRetrievalType, logRetrievalScope, 1186 logFilter, ps); 1187 } 1188 finally { 1189 ps.close(); 1190 } 1191 } 1192 else { 1193 if (!options.contains(ACTION_OPTION) && !options.contains(DATE_OPTION)) { 1194 wc.getJobLog(commandLine.getOptionValue(LOG_OPTION), null, null, logFilter, ps); 1195 } 1196 else { 1197 throw new OozieCLIException("Invalid options provided for log retrieval. " + ACTION_OPTION 1198 + " and " + DATE_OPTION + " are valid only for coordinator job log retrieval"); 1199 } 1200 } 1201 } 1202 else if (options.contains(DEFINITION_OPTION)) { 1203 System.out.println(wc.getJobDefinition(commandLine.getOptionValue(DEFINITION_OPTION))); 1204 } 1205 else if (options.contains(CONFIG_CONTENT_OPTION)) { 1206 if (commandLine.getOptionValue(CONFIG_CONTENT_OPTION).endsWith("-C")) { 1207 System.out.println(wc.getCoordJobInfo(commandLine.getOptionValue(CONFIG_CONTENT_OPTION)).getConf()); 1208 } 1209 else if (commandLine.getOptionValue(CONFIG_CONTENT_OPTION).endsWith("-W")) { 1210 System.out.println(wc.getJobInfo(commandLine.getOptionValue(CONFIG_CONTENT_OPTION)).getConf()); 1211 } 1212 else if (commandLine.getOptionValue(CONFIG_CONTENT_OPTION).endsWith("-B")) { 1213 System.out 1214 .println(wc.getBundleJobInfo(commandLine.getOptionValue(CONFIG_CONTENT_OPTION)).getConf()); 1215 } 1216 else { 1217 System.out.println("ERROR: job id [" + commandLine.getOptionValue(CONFIG_CONTENT_OPTION) 1218 + "] doesn't end with either C or W or B"); 1219 } 1220 } 1221 else if (options.contains(UPDATE_OPTION)) { 1222 String coordJobId = commandLine.getOptionValue(UPDATE_OPTION); 1223 Properties conf = null; 1224 1225 String dryrun = ""; 1226 String showdiff = ""; 1227 1228 if (commandLine.getOptionValue(CONFIG_OPTION) != null) { 1229 conf = getConfiguration(wc, commandLine); 1230 } 1231 if (options.contains(DRYRUN_OPTION)) { 1232 dryrun = "true"; 1233 } 1234 if (commandLine.getOptionValue(SHOWDIFF_OPTION) != null) { 1235 showdiff = commandLine.getOptionValue(SHOWDIFF_OPTION); 1236 } 1237 if (conf == null) { 1238 System.out.println(wc.updateCoord(coordJobId, dryrun, showdiff)); 1239 } 1240 else { 1241 System.out.println(wc.updateCoord(coordJobId, conf, dryrun, showdiff)); 1242 } 1243 } 1244 else if (options.contains(POLL_OPTION)) { 1245 String jobId = commandLine.getOptionValue(POLL_OPTION); 1246 int timeout = 30; 1247 int interval = 5; 1248 String timeoutS = commandLine.getOptionValue(TIMEOUT_OPTION); 1249 if (timeoutS != null) { 1250 timeout = Integer.parseInt(timeoutS); 1251 } 1252 String intervalS = commandLine.getOptionValue(INTERVAL_OPTION); 1253 if (intervalS != null) { 1254 interval = Integer.parseInt(intervalS); 1255 } 1256 boolean verbose = commandLine.hasOption(VERBOSE_OPTION); 1257 wc.pollJob(jobId, timeout, interval, verbose); 1258 } 1259 } 1260 catch (OozieClientException ex) { 1261 throw new OozieCLIException(ex.toString(), ex); 1262 } 1263 } 1264 1265 @VisibleForTesting 1266 void printCoordJob(CoordinatorJob coordJob, String timeZoneId, boolean verbose) { 1267 System.out.println("Job ID : " + coordJob.getId()); 1268 1269 System.out.println(RULER); 1270 1271 List<CoordinatorAction> actions = coordJob.getActions(); 1272 System.out.println("Job Name : " + maskIfNull(coordJob.getAppName())); 1273 System.out.println("App Path : " + maskIfNull(coordJob.getAppPath())); 1274 System.out.println("Status : " + coordJob.getStatus()); 1275 System.out.println("Start Time : " + maskDate(coordJob.getStartTime(), timeZoneId, false)); 1276 System.out.println("End Time : " + maskDate(coordJob.getEndTime(), timeZoneId, false)); 1277 System.out.println("Pause Time : " + maskDate(coordJob.getPauseTime(), timeZoneId, false)); 1278 System.out.println("Concurrency : " + coordJob.getConcurrency()); 1279 System.out.println(RULER); 1280 1281 if (verbose) { 1282 System.out.println("ID" + VERBOSE_DELIMITER + "Action Number" + VERBOSE_DELIMITER + "Console URL" 1283 + VERBOSE_DELIMITER + "Error Code" + VERBOSE_DELIMITER + "Error Message" + VERBOSE_DELIMITER 1284 + "External ID" + VERBOSE_DELIMITER + "External Status" + VERBOSE_DELIMITER + "Job ID" 1285 + VERBOSE_DELIMITER + "Tracker URI" + VERBOSE_DELIMITER + "Created" + VERBOSE_DELIMITER 1286 + "Nominal Time" + VERBOSE_DELIMITER + "Status" + VERBOSE_DELIMITER + "Last Modified" 1287 + VERBOSE_DELIMITER + "Missing Dependencies"); 1288 System.out.println(RULER); 1289 1290 for (CoordinatorAction action : actions) { 1291 System.out.println(maskIfNull(action.getId()) + VERBOSE_DELIMITER + action.getActionNumber() 1292 + VERBOSE_DELIMITER + maskIfNull(action.getConsoleUrl()) + VERBOSE_DELIMITER 1293 + maskIfNull(action.getErrorCode()) + VERBOSE_DELIMITER + maskIfNull(action.getErrorMessage()) 1294 + VERBOSE_DELIMITER + maskIfNull(action.getExternalId()) + VERBOSE_DELIMITER 1295 + maskIfNull(action.getExternalStatus()) + VERBOSE_DELIMITER + maskIfNull(action.getJobId()) 1296 + VERBOSE_DELIMITER + maskIfNull(action.getTrackerUri()) + VERBOSE_DELIMITER 1297 + maskDate(action.getCreatedTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1298 + maskDate(action.getNominalTime(), timeZoneId, verbose) + action.getStatus() + VERBOSE_DELIMITER 1299 + maskDate(action.getLastModifiedTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1300 + maskIfNull(getFirstMissingDependencies(action))); 1301 1302 System.out.println(RULER); 1303 } 1304 } 1305 else { 1306 System.out.println(String.format(COORD_ACTION_FORMATTER, "ID", "Status", "Ext ID", "Err Code", "Created", 1307 "Nominal Time", "Last Mod")); 1308 1309 for (CoordinatorAction action : actions) { 1310 System.out.println(String.format(COORD_ACTION_FORMATTER, maskIfNull(action.getId()), 1311 action.getStatus(), maskIfNull(action.getExternalId()), maskIfNull(action.getErrorCode()), 1312 maskDate(action.getCreatedTime(), timeZoneId, verbose), maskDate(action.getNominalTime(), timeZoneId, verbose), 1313 maskDate(action.getLastModifiedTime(), timeZoneId, verbose))); 1314 1315 System.out.println(RULER); 1316 } 1317 } 1318 } 1319 1320 @VisibleForTesting 1321 void printBundleJob(BundleJob bundleJob, String timeZoneId, boolean verbose) { 1322 System.out.println("Job ID : " + bundleJob.getId()); 1323 1324 System.out.println(RULER); 1325 1326 List<CoordinatorJob> coordinators = bundleJob.getCoordinators(); 1327 System.out.println("Job Name : " + maskIfNull(bundleJob.getAppName())); 1328 System.out.println("App Path : " + maskIfNull(bundleJob.getAppPath())); 1329 System.out.println("Status : " + bundleJob.getStatus()); 1330 System.out.println("Kickoff time : " + bundleJob.getKickoffTime()); 1331 System.out.println(RULER); 1332 1333 System.out.println(String.format(BUNDLE_COORD_JOBS_FORMATTER, "Job ID", "Status", "Freq", "Unit", "Started", 1334 "Next Materialized")); 1335 System.out.println(RULER); 1336 1337 for (CoordinatorJob job : coordinators) { 1338 System.out.println(String.format(BUNDLE_COORD_JOBS_FORMATTER, maskIfNull(job.getId()), job.getStatus(), 1339 job.getFrequency(), job.getTimeUnit(), maskDate(job.getStartTime(), timeZoneId, verbose), 1340 maskDate(job.getNextMaterializedTime(), timeZoneId, verbose))); 1341 1342 System.out.println(RULER); 1343 } 1344 } 1345 1346 @VisibleForTesting 1347 void printCoordAction(CoordinatorAction coordAction, String timeZoneId) { 1348 System.out.println("ID : " + maskIfNull(coordAction.getId())); 1349 1350 System.out.println(RULER); 1351 1352 System.out.println("Action Number : " + coordAction.getActionNumber()); 1353 System.out.println("Console URL : " + maskIfNull(coordAction.getConsoleUrl())); 1354 System.out.println("Error Code : " + maskIfNull(coordAction.getErrorCode())); 1355 System.out.println("Error Message : " + maskIfNull(coordAction.getErrorMessage())); 1356 System.out.println("External ID : " + maskIfNull(coordAction.getExternalId())); 1357 System.out.println("External Status : " + maskIfNull(coordAction.getExternalStatus())); 1358 System.out.println("Job ID : " + maskIfNull(coordAction.getJobId())); 1359 System.out.println("Tracker URI : " + maskIfNull(coordAction.getTrackerUri())); 1360 System.out.println("Created : " + maskDate(coordAction.getCreatedTime(), timeZoneId, false)); 1361 System.out.println("Nominal Time : " + maskDate(coordAction.getNominalTime(), timeZoneId, false)); 1362 System.out.println("Status : " + coordAction.getStatus()); 1363 System.out.println("Last Modified : " + maskDate(coordAction.getLastModifiedTime(), timeZoneId, false)); 1364 System.out.println("First Missing Dependency : " + maskIfNull(getFirstMissingDependencies(coordAction))); 1365 1366 System.out.println(RULER); 1367 } 1368 1369 private void printCoordActions(List<CoordinatorAction> actions) { 1370 if (actions != null && actions.size() > 0) { 1371 System.out.println("Action ID" + VERBOSE_DELIMITER + "Nominal Time"); 1372 System.out.println(RULER); 1373 for (CoordinatorAction action : actions) { 1374 System.out.println(maskIfNull(action.getId()) + VERBOSE_DELIMITER 1375 + maskDate(action.getNominalTime(), null,false)); 1376 } 1377 } 1378 else { 1379 System.out.println("No Actions match your criteria!"); 1380 } 1381 } 1382 1383 private void printCoordActionsStatus(List<CoordinatorAction> actions) { 1384 if (actions != null && actions.size() > 0) { 1385 System.out.println("Action ID" + VERBOSE_DELIMITER + "Nominal Time" + VERBOSE_DELIMITER + "Status"); 1386 System.out.println(RULER); 1387 for (CoordinatorAction action : actions) { 1388 System.out.println(maskIfNull(action.getId()) + VERBOSE_DELIMITER 1389 + maskDate(action.getNominalTime(), null, false) + VERBOSE_DELIMITER 1390 + maskIfNull(action.getStatus().name())); 1391 } 1392 } 1393 } 1394 1395 @VisibleForTesting 1396 void printWorkflowAction(WorkflowAction action, String timeZoneId, boolean verbose) { 1397 1398 System.out.println("ID : " + maskIfNull(action.getId())); 1399 1400 System.out.println(RULER); 1401 1402 System.out.println("Console URL : " + maskIfNull(action.getConsoleUrl())); 1403 System.out.println("Error Code : " + maskIfNull(action.getErrorCode())); 1404 System.out.println("Error Message : " + maskIfNull(action.getErrorMessage())); 1405 System.out.println("External ID : " + maskIfNull(action.getExternalId())); 1406 System.out.println("External Status : " + maskIfNull(action.getExternalStatus())); 1407 System.out.println("Name : " + maskIfNull(action.getName())); 1408 System.out.println("Retries : " + action.getRetries()); 1409 System.out.println("Tracker URI : " + maskIfNull(action.getTrackerUri())); 1410 System.out.println("Type : " + maskIfNull(action.getType())); 1411 System.out.println("Started : " + maskDate(action.getStartTime(), timeZoneId, verbose)); 1412 System.out.println("Status : " + action.getStatus()); 1413 System.out.println("Ended : " + maskDate(action.getEndTime(), timeZoneId, verbose)); 1414 1415 if (verbose) { 1416 System.out.println("External Stats : " + action.getStats()); 1417 System.out.println("External ChildIDs : " + action.getExternalChildIDs()); 1418 } 1419 1420 System.out.println(RULER); 1421 } 1422 1423 private static final String WORKFLOW_JOBS_FORMATTER = "%-41s%-13s%-10s%-10s%-10s%-24s%-24s"; 1424 private static final String COORD_JOBS_FORMATTER = "%-41s%-15s%-10s%-5s%-13s%-24s%-24s"; 1425 private static final String BUNDLE_JOBS_FORMATTER = "%-41s%-15s%-10s%-20s%-20s%-13s%-13s"; 1426 private static final String BUNDLE_COORD_JOBS_FORMATTER = "%-41s%-15s%-5s%-13s%-24s%-24s"; 1427 1428 private static final String WORKFLOW_ACTION_FORMATTER = "%-78s%-10s%-23s%-11s%-10s"; 1429 private static final String COORD_ACTION_FORMATTER = "%-43s%-10s%-37s%-10s%-21s%-21s"; 1430 private static final String BULK_RESPONSE_FORMATTER = "%-13s%-38s%-13s%-41s%-10s%-38s%-21s%-38s"; 1431 1432 @VisibleForTesting 1433 void printJob(WorkflowJob job, String timeZoneId, boolean verbose) throws IOException { 1434 System.out.println("Job ID : " + maskIfNull(job.getId())); 1435 1436 System.out.println(RULER); 1437 1438 System.out.println("Workflow Name : " + maskIfNull(job.getAppName())); 1439 System.out.println("App Path : " + maskIfNull(job.getAppPath())); 1440 System.out.println("Status : " + job.getStatus()); 1441 System.out.println("Run : " + job.getRun()); 1442 System.out.println("User : " + maskIfNull(job.getUser())); 1443 System.out.println("Group : " + maskIfNull(job.getGroup())); 1444 System.out.println("Created : " + maskDate(job.getCreatedTime(), timeZoneId, verbose)); 1445 System.out.println("Started : " + maskDate(job.getStartTime(), timeZoneId, verbose)); 1446 System.out.println("Last Modified : " + maskDate(job.getLastModifiedTime(), timeZoneId, verbose)); 1447 System.out.println("Ended : " + maskDate(job.getEndTime(), timeZoneId, verbose)); 1448 System.out.println("CoordAction ID: " + maskIfNull(job.getParentId())); 1449 1450 List<WorkflowAction> actions = job.getActions(); 1451 1452 if (actions != null && actions.size() > 0) { 1453 System.out.println(); 1454 System.out.println("Actions"); 1455 System.out.println(RULER); 1456 1457 if (verbose) { 1458 System.out.println("ID" + VERBOSE_DELIMITER + "Console URL" + VERBOSE_DELIMITER + "Error Code" 1459 + VERBOSE_DELIMITER + "Error Message" + VERBOSE_DELIMITER + "External ID" + VERBOSE_DELIMITER 1460 + "External Status" + VERBOSE_DELIMITER + "Name" + VERBOSE_DELIMITER + "Retries" 1461 + VERBOSE_DELIMITER + "Tracker URI" + VERBOSE_DELIMITER + "Type" + VERBOSE_DELIMITER 1462 + "Started" + VERBOSE_DELIMITER + "Status" + VERBOSE_DELIMITER + "Ended"); 1463 System.out.println(RULER); 1464 1465 for (WorkflowAction action : job.getActions()) { 1466 System.out.println(maskIfNull(action.getId()) + VERBOSE_DELIMITER 1467 + maskIfNull(action.getConsoleUrl()) + VERBOSE_DELIMITER 1468 + maskIfNull(action.getErrorCode()) + VERBOSE_DELIMITER 1469 + maskIfNull(action.getErrorMessage()) + VERBOSE_DELIMITER 1470 + maskIfNull(action.getExternalId()) + VERBOSE_DELIMITER 1471 + maskIfNull(action.getExternalStatus()) + VERBOSE_DELIMITER + maskIfNull(action.getName()) 1472 + VERBOSE_DELIMITER + action.getRetries() + VERBOSE_DELIMITER 1473 + maskIfNull(action.getTrackerUri()) + VERBOSE_DELIMITER + maskIfNull(action.getType()) 1474 + VERBOSE_DELIMITER + maskDate(action.getStartTime(), timeZoneId, verbose) 1475 + VERBOSE_DELIMITER + action.getStatus() + VERBOSE_DELIMITER 1476 + maskDate(action.getEndTime(), timeZoneId, verbose)); 1477 1478 System.out.println(RULER); 1479 } 1480 } 1481 else { 1482 System.out.println(String.format(WORKFLOW_ACTION_FORMATTER, "ID", "Status", "Ext ID", "Ext Status", 1483 "Err Code")); 1484 1485 System.out.println(RULER); 1486 1487 for (WorkflowAction action : job.getActions()) { 1488 System.out.println(String.format(WORKFLOW_ACTION_FORMATTER, maskIfNull(action.getId()), action 1489 .getStatus(), maskIfNull(action.getExternalId()), maskIfNull(action.getExternalStatus()), 1490 maskIfNull(action.getErrorCode()))); 1491 1492 System.out.println(RULER); 1493 } 1494 } 1495 } 1496 else { 1497 System.out.println(RULER); 1498 } 1499 1500 System.out.println(); 1501 } 1502 1503 private void jobsCommand(CommandLine commandLine) throws IOException, OozieCLIException { 1504 XOozieClient wc = createXOozieClient(commandLine); 1505 1506 List<String> options = new ArrayList<String>(); 1507 for (Option option : commandLine.getOptions()) { 1508 options.add(option.getOpt()); 1509 } 1510 1511 String filter = commandLine.getOptionValue(FILTER_OPTION); 1512 String s = commandLine.getOptionValue(OFFSET_OPTION); 1513 int start = Integer.parseInt((s != null) ? s : "0"); 1514 s = commandLine.getOptionValue(LEN_OPTION); 1515 String jobtype = commandLine.getOptionValue(JOBTYPE_OPTION); 1516 String timeZoneId = getTimeZoneId(commandLine); 1517 jobtype = (jobtype != null) ? jobtype : "wf"; 1518 int len = Integer.parseInt((s != null) ? s : "0"); 1519 String bulkFilterString = commandLine.getOptionValue(BULK_OPTION); 1520 1521 try { 1522 if (options.contains(KILL_OPTION)) { 1523 printBulkModifiedJobs(wc.killJobs(filter, jobtype, start, len), timeZoneId, "killed"); 1524 } 1525 else if (options.contains(SUSPEND_OPTION)) { 1526 printBulkModifiedJobs(wc.suspendJobs(filter, jobtype, start, len), timeZoneId, "suspended"); 1527 } 1528 else if (options.contains(RESUME_OPTION)) { 1529 printBulkModifiedJobs(wc.resumeJobs(filter, jobtype, start, len), timeZoneId, "resumed"); 1530 } 1531 else if (bulkFilterString != null) { 1532 printBulkJobs(wc.getBulkInfo(bulkFilterString, start, len), timeZoneId, commandLine.hasOption(VERBOSE_OPTION)); 1533 } 1534 else if (jobtype.toLowerCase().contains("wf")) { 1535 printJobs(wc.getJobsInfo(filter, start, len), timeZoneId, commandLine.hasOption(VERBOSE_OPTION)); 1536 } 1537 else if (jobtype.toLowerCase().startsWith("coord")) { 1538 printCoordJobs(wc.getCoordJobsInfo(filter, start, len), timeZoneId, commandLine.hasOption(VERBOSE_OPTION)); 1539 } 1540 else if (jobtype.toLowerCase().startsWith("bundle")) { 1541 printBundleJobs(wc.getBundleJobsInfo(filter, start, len), timeZoneId, commandLine.hasOption(VERBOSE_OPTION)); 1542 } 1543 1544 } 1545 catch (OozieClientException ex) { 1546 throw new OozieCLIException(ex.toString(), ex); 1547 } 1548 } 1549 1550 @VisibleForTesting 1551 void printBulkModifiedJobs(JSONObject json, String timeZoneId, String action) throws IOException { 1552 if (json.containsKey(JsonTags.WORKFLOWS_JOBS)) { 1553 JSONArray workflows = (JSONArray) json.get(JsonTags.WORKFLOWS_JOBS); 1554 if (workflows == null) { 1555 workflows = new JSONArray(); 1556 } 1557 List<WorkflowJob> wfs = JsonToBean.createWorkflowJobList(workflows); 1558 if (wfs.isEmpty()) { 1559 System.out.println("bulk modify command did not modify any jobs"); 1560 } 1561 else { 1562 System.out.println("the following jobs have been " + action); 1563 printJobs(wfs, timeZoneId, false); 1564 } 1565 } 1566 else if (json.containsKey(JsonTags.COORDINATOR_JOBS)) { 1567 JSONArray coordinators = (JSONArray) json.get(JsonTags.COORDINATOR_JOBS); 1568 if (coordinators == null) { 1569 coordinators = new JSONArray(); 1570 } 1571 List<CoordinatorJob> coords = JsonToBean.createCoordinatorJobList(coordinators); 1572 if (coords.isEmpty()) { 1573 System.out.println("bulk modify command did not modify any jobs"); 1574 } 1575 else { 1576 System.out.println("the following jobs have been " + action); 1577 printCoordJobs(coords, timeZoneId, false); 1578 } 1579 } 1580 else { 1581 JSONArray bundles = (JSONArray) json.get(JsonTags.BUNDLE_JOBS); 1582 if (bundles == null) { 1583 bundles = new JSONArray(); 1584 } 1585 List<BundleJob> bundleJobs = JsonToBean.createBundleJobList(bundles); 1586 if (bundleJobs.isEmpty()) { 1587 System.out.println("bulk modify command did not modify any jobs"); 1588 } 1589 else { 1590 System.out.println("the following jobs have been " + action); 1591 printBundleJobs(bundleJobs, timeZoneId, false); 1592 } 1593 } 1594 } 1595 1596 @VisibleForTesting 1597 void printCoordJobs(List<CoordinatorJob> jobs, String timeZoneId, boolean verbose) throws IOException { 1598 if (jobs != null && jobs.size() > 0) { 1599 if (verbose) { 1600 System.out.println("Job ID" + VERBOSE_DELIMITER + "App Name" + VERBOSE_DELIMITER + "App Path" 1601 + VERBOSE_DELIMITER + "Console URL" + VERBOSE_DELIMITER + "User" + VERBOSE_DELIMITER + "Group" 1602 + VERBOSE_DELIMITER + "Concurrency" + VERBOSE_DELIMITER + "Frequency" + VERBOSE_DELIMITER 1603 + "Time Unit" + VERBOSE_DELIMITER + "Time Zone" + VERBOSE_DELIMITER + "Time Out" 1604 + VERBOSE_DELIMITER + "Started" + VERBOSE_DELIMITER + "Next Materialize" + VERBOSE_DELIMITER 1605 + "Status" + VERBOSE_DELIMITER + "Last Action" + VERBOSE_DELIMITER + "Ended"); 1606 System.out.println(RULER); 1607 1608 for (CoordinatorJob job : jobs) { 1609 System.out.println(maskIfNull(job.getId()) + VERBOSE_DELIMITER + maskIfNull(job.getAppName()) 1610 + VERBOSE_DELIMITER + maskIfNull(job.getAppPath()) + VERBOSE_DELIMITER 1611 + maskIfNull(job.getConsoleUrl()) + VERBOSE_DELIMITER + maskIfNull(job.getUser()) 1612 + VERBOSE_DELIMITER + maskIfNull(job.getGroup()) + VERBOSE_DELIMITER + job.getConcurrency() 1613 + VERBOSE_DELIMITER + job.getFrequency() + VERBOSE_DELIMITER + job.getTimeUnit() 1614 + VERBOSE_DELIMITER + maskIfNull(job.getTimeZone()) + VERBOSE_DELIMITER + job.getTimeout() 1615 + VERBOSE_DELIMITER + maskDate(job.getStartTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1616 + maskDate(job.getNextMaterializedTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1617 + job.getStatus() + VERBOSE_DELIMITER 1618 + maskDate(job.getLastActionTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1619 + maskDate(job.getEndTime(), timeZoneId, verbose)); 1620 1621 System.out.println(RULER); 1622 } 1623 } 1624 else { 1625 System.out.println(String.format(COORD_JOBS_FORMATTER, "Job ID", "App Name", "Status", "Freq", "Unit", 1626 "Started", "Next Materialized")); 1627 System.out.println(RULER); 1628 1629 for (CoordinatorJob job : jobs) { 1630 System.out.println(String.format(COORD_JOBS_FORMATTER, maskIfNull(job.getId()), maskIfNull(job 1631 .getAppName()), job.getStatus(), job.getFrequency(), job.getTimeUnit(), maskDate(job 1632 .getStartTime(), timeZoneId, verbose), maskDate(job.getNextMaterializedTime(), timeZoneId, verbose))); 1633 1634 System.out.println(RULER); 1635 } 1636 } 1637 } 1638 else { 1639 System.out.println("No Jobs match your criteria!"); 1640 } 1641 } 1642 1643 @VisibleForTesting 1644 void printBulkJobs(List<BulkResponse> jobs, String timeZoneId, boolean verbose) throws IOException { 1645 if (jobs != null && jobs.size() > 0) { 1646 for (BulkResponse response : jobs) { 1647 BundleJob bundle = response.getBundle(); 1648 CoordinatorJob coord = response.getCoordinator(); 1649 CoordinatorAction action = response.getAction(); 1650 if (verbose) { 1651 System.out.println(); 1652 System.out.println("Bundle Name : " + maskIfNull(bundle.getAppName())); 1653 1654 System.out.println(RULER); 1655 1656 System.out.println("Bundle ID : " + maskIfNull(bundle.getId())); 1657 System.out.println("Coordinator Name : " + maskIfNull(coord.getAppName())); 1658 System.out.println("Coord Action ID : " + maskIfNull(action.getId())); 1659 System.out.println("Action Status : " + action.getStatus()); 1660 System.out.println("External ID : " + maskIfNull(action.getExternalId())); 1661 System.out.println("Created Time : " + maskDate(action.getCreatedTime(), timeZoneId, false)); 1662 System.out.println("User : " + maskIfNull(bundle.getUser())); 1663 System.out.println("Error Message : " + maskIfNull(action.getErrorMessage())); 1664 System.out.println(RULER); 1665 } 1666 else { 1667 System.out.println(String.format(BULK_RESPONSE_FORMATTER, "Bundle Name", "Bundle ID", "Coord Name", 1668 "Coord Action ID", "Status", "External ID", "Created Time", "Error Message")); 1669 System.out.println(RULER); 1670 System.out 1671 .println(String.format(BULK_RESPONSE_FORMATTER, maskIfNull(bundle.getAppName()), 1672 maskIfNull(bundle.getId()), maskIfNull(coord.getAppName()), 1673 maskIfNull(action.getId()), action.getStatus(), maskIfNull(action.getExternalId()), 1674 maskDate(action.getCreatedTime(), timeZoneId, false), 1675 maskIfNull(action.getErrorMessage()))); 1676 System.out.println(RULER); 1677 } 1678 } 1679 } 1680 else { 1681 System.out.println("Bulk request criteria did not match any coordinator actions"); 1682 } 1683 } 1684 1685 @VisibleForTesting 1686 void printBundleJobs(List<BundleJob> jobs, String timeZoneId, boolean verbose) throws IOException { 1687 if (jobs != null && jobs.size() > 0) { 1688 if (verbose) { 1689 System.out.println("Job ID" + VERBOSE_DELIMITER + "Bundle Name" + VERBOSE_DELIMITER + "Bundle Path" 1690 + VERBOSE_DELIMITER + "User" + VERBOSE_DELIMITER + "Group" + VERBOSE_DELIMITER + "Status" 1691 + VERBOSE_DELIMITER + "Kickoff" + VERBOSE_DELIMITER + "Pause" + VERBOSE_DELIMITER + "Created" 1692 + VERBOSE_DELIMITER + "Console URL"); 1693 System.out.println(RULER); 1694 1695 for (BundleJob job : jobs) { 1696 System.out.println(maskIfNull(job.getId()) + VERBOSE_DELIMITER + maskIfNull(job.getAppName()) 1697 + VERBOSE_DELIMITER + maskIfNull(job.getAppPath()) + VERBOSE_DELIMITER 1698 + maskIfNull(job.getUser()) + VERBOSE_DELIMITER + maskIfNull(job.getGroup()) 1699 + VERBOSE_DELIMITER + job.getStatus() + VERBOSE_DELIMITER 1700 + maskDate(job.getKickoffTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1701 + maskDate(job.getPauseTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1702 + maskDate(job.getCreatedTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1703 + maskIfNull(job.getConsoleUrl())); 1704 1705 System.out.println(RULER); 1706 } 1707 } 1708 else { 1709 System.out.println(String.format(BUNDLE_JOBS_FORMATTER, "Job ID", "Bundle Name", "Status", "Kickoff", 1710 "Created", "User", "Group")); 1711 System.out.println(RULER); 1712 1713 for (BundleJob job : jobs) { 1714 System.out.println(String.format(BUNDLE_JOBS_FORMATTER, maskIfNull(job.getId()), 1715 maskIfNull(job.getAppName()), job.getStatus(), 1716 maskDate(job.getKickoffTime(), timeZoneId, verbose), 1717 maskDate(job.getCreatedTime(), timeZoneId, verbose), maskIfNull(job.getUser()), 1718 maskIfNull(job.getGroup()))); 1719 System.out.println(RULER); 1720 } 1721 } 1722 } 1723 else { 1724 System.out.println("No Jobs match your criteria!"); 1725 } 1726 } 1727 1728 private void slaCommand(CommandLine commandLine) throws IOException, OozieCLIException { 1729 XOozieClient wc = createXOozieClient(commandLine); 1730 List<String> options = new ArrayList<String>(); 1731 for (Option option : commandLine.getOptions()) { 1732 options.add(option.getOpt()); 1733 } 1734 1735 String s = commandLine.getOptionValue(OFFSET_OPTION); 1736 int start = Integer.parseInt((s != null) ? s : "0"); 1737 s = commandLine.getOptionValue(LEN_OPTION); 1738 int len = Integer.parseInt((s != null) ? s : "100"); 1739 String filter = commandLine.getOptionValue(FILTER_OPTION); 1740 1741 try { 1742 wc.getSlaInfo(start, len, filter); 1743 } 1744 catch (OozieClientException ex) { 1745 throw new OozieCLIException(ex.toString(), ex); 1746 } 1747 } 1748 1749 private void adminCommand(CommandLine commandLine) throws OozieCLIException { 1750 XOozieClient wc = createXOozieClient(commandLine); 1751 1752 List<String> options = new ArrayList<String>(); 1753 for (Option option : commandLine.getOptions()) { 1754 options.add(option.getOpt()); 1755 } 1756 1757 try { 1758 SYSTEM_MODE status = SYSTEM_MODE.NORMAL; 1759 if (options.contains(VERSION_OPTION)) { 1760 System.out.println("Oozie server build version: " + wc.getServerBuildVersion()); 1761 } 1762 else if (options.contains(SYSTEM_MODE_OPTION)) { 1763 String systemModeOption = commandLine.getOptionValue(SYSTEM_MODE_OPTION).toUpperCase(); 1764 try { 1765 status = SYSTEM_MODE.valueOf(systemModeOption); 1766 } 1767 catch (Exception e) { 1768 throw new OozieCLIException("Invalid input provided for option: " + SYSTEM_MODE_OPTION 1769 + " value given :" + systemModeOption 1770 + " Expected values are: NORMAL/NOWEBSERVICE/SAFEMODE "); 1771 } 1772 wc.setSystemMode(status); 1773 System.out.println("System mode: " + status); 1774 } 1775 else if (options.contains(STATUS_OPTION)) { 1776 status = wc.getSystemMode(); 1777 System.out.println("System mode: " + status); 1778 } 1779 1780 else if (options.contains(UPDATE_SHARELIB_OPTION)) { 1781 System.out.println(wc.updateShareLib()); 1782 } 1783 1784 else if (options.contains(LIST_SHARELIB_LIB_OPTION)) { 1785 String sharelibKey = null; 1786 if (commandLine.getArgList().size() > 0) { 1787 sharelibKey = (String) commandLine.getArgList().get(0); 1788 } 1789 System.out.println(wc.listShareLib(sharelibKey)); 1790 } 1791 1792 else if (options.contains(QUEUE_DUMP_OPTION)) { 1793 1794 List<String> list = wc.getQueueDump(); 1795 if (list != null && list.size() != 0) { 1796 for (String str : list) { 1797 System.out.println(str); 1798 } 1799 } 1800 else { 1801 System.out.println("QueueDump is null!"); 1802 } 1803 } 1804 else if (options.contains(AVAILABLE_SERVERS_OPTION)) { 1805 Map<String, String> availableOozieServers = new TreeMap<String, String>(wc.getAvailableOozieServers()); 1806 for (Map.Entry<String, String> ent : availableOozieServers.entrySet()) { 1807 System.out.println(ent.getKey() + " : " + ent.getValue()); 1808 } 1809 } else if (options.contains(SERVER_CONFIGURATION_OPTION)) { 1810 Map<String, String> serverConfig = new TreeMap<String, String>(wc.getServerConfiguration()); 1811 for (Map.Entry<String, String> ent : serverConfig.entrySet()) { 1812 System.out.println(ent.getKey() + " : " + ent.getValue()); 1813 } 1814 } else if (options.contains(SERVER_OS_ENV_OPTION)) { 1815 Map<String, String> osEnv = new TreeMap<String, String>(wc.getOSEnv()); 1816 for (Map.Entry<String, String> ent : osEnv.entrySet()) { 1817 System.out.println(ent.getKey() + " : " + ent.getValue()); 1818 } 1819 } else if (options.contains(SERVER_JAVA_SYSTEM_PROPERTIES_OPTION)) { 1820 Map<String, String> javaSysProps = new TreeMap<String, String>(wc.getJavaSystemProperties()); 1821 for (Map.Entry<String, String> ent : javaSysProps.entrySet()) { 1822 System.out.println(ent.getKey() + " : " + ent.getValue()); 1823 } 1824 } else if (options.contains(METRICS_OPTION)) { 1825 OozieClient.Metrics metrics = wc.getMetrics(); 1826 if (metrics == null) { 1827 System.out.println("Metrics are unavailable. Try Instrumentation (-" + INSTRUMENTATION_OPTION + ") instead"); 1828 } else { 1829 printMetrics(metrics); 1830 } 1831 } else if (options.contains(INSTRUMENTATION_OPTION)) { 1832 OozieClient.Instrumentation instrumentation = wc.getInstrumentation(); 1833 if (instrumentation == null) { 1834 System.out.println("Instrumentation is unavailable. Try Metrics (-" + METRICS_OPTION + ") instead"); 1835 } else { 1836 printInstrumentation(instrumentation); 1837 } 1838 } 1839 } 1840 catch (OozieClientException ex) { 1841 throw new OozieCLIException(ex.toString(), ex); 1842 } 1843 } 1844 1845 private void versionCommand() throws OozieCLIException { 1846 System.out.println("Oozie client build version: " 1847 + BuildInfo.getBuildInfo().getProperty(BuildInfo.BUILD_VERSION)); 1848 } 1849 1850 @VisibleForTesting 1851 void printJobs(List<WorkflowJob> jobs, String timeZoneId, boolean verbose) throws IOException { 1852 if (jobs != null && jobs.size() > 0) { 1853 if (verbose) { 1854 System.out.println("Job ID" + VERBOSE_DELIMITER + "App Name" + VERBOSE_DELIMITER + "App Path" 1855 + VERBOSE_DELIMITER + "Console URL" + VERBOSE_DELIMITER + "User" + VERBOSE_DELIMITER + "Group" 1856 + VERBOSE_DELIMITER + "Run" + VERBOSE_DELIMITER + "Created" + VERBOSE_DELIMITER + "Started" 1857 + VERBOSE_DELIMITER + "Status" + VERBOSE_DELIMITER + "Last Modified" + VERBOSE_DELIMITER 1858 + "Ended"); 1859 System.out.println(RULER); 1860 1861 for (WorkflowJob job : jobs) { 1862 System.out.println(maskIfNull(job.getId()) + VERBOSE_DELIMITER + maskIfNull(job.getAppName()) 1863 + VERBOSE_DELIMITER + maskIfNull(job.getAppPath()) + VERBOSE_DELIMITER 1864 + maskIfNull(job.getConsoleUrl()) + VERBOSE_DELIMITER + maskIfNull(job.getUser()) 1865 + VERBOSE_DELIMITER + maskIfNull(job.getGroup()) + VERBOSE_DELIMITER + job.getRun() 1866 + VERBOSE_DELIMITER + maskDate(job.getCreatedTime(), timeZoneId, verbose) 1867 + VERBOSE_DELIMITER + maskDate(job.getStartTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1868 + job.getStatus() + VERBOSE_DELIMITER 1869 + maskDate(job.getLastModifiedTime(), timeZoneId, verbose) + VERBOSE_DELIMITER 1870 + maskDate(job.getEndTime(), timeZoneId, verbose)); 1871 1872 System.out.println(RULER); 1873 } 1874 } 1875 else { 1876 System.out.println(String.format(WORKFLOW_JOBS_FORMATTER, "Job ID", "App Name", "Status", "User", 1877 "Group", "Started", "Ended")); 1878 System.out.println(RULER); 1879 1880 for (WorkflowJob job : jobs) { 1881 System.out.println(String.format(WORKFLOW_JOBS_FORMATTER, maskIfNull(job.getId()), 1882 maskIfNull(job.getAppName()), job.getStatus(), maskIfNull(job.getUser()), 1883 maskIfNull(job.getGroup()), maskDate(job.getStartTime(), timeZoneId, verbose), 1884 maskDate(job.getEndTime(), timeZoneId, verbose))); 1885 1886 System.out.println(RULER); 1887 } 1888 } 1889 } 1890 else { 1891 System.out.println("No Jobs match your criteria!"); 1892 } 1893 } 1894 1895 void printWfsForCoordAction(List<WorkflowJob> jobs, String timeZoneId) throws IOException { 1896 if (jobs != null && jobs.size() > 0) { 1897 System.out.println(String.format("%-41s%-10s%-24s%-24s", "Job ID", "Status", "Started", "Ended")); 1898 System.out.println(RULER); 1899 1900 for (WorkflowJob job : jobs) { 1901 System.out 1902 .println(String.format("%-41s%-10s%-24s%-24s", maskIfNull(job.getId()), job.getStatus(), 1903 maskDate(job.getStartTime(), timeZoneId, false), 1904 maskDate(job.getEndTime(), timeZoneId, false))); 1905 System.out.println(RULER); 1906 } 1907 } 1908 } 1909 1910 private String maskIfNull(String value) { 1911 if (value != null && value.length() > 0) { 1912 return value; 1913 } 1914 return "-"; 1915 } 1916 1917 private String maskDate(Date date, String timeZoneId, boolean verbose) { 1918 if (date == null) { 1919 return "-"; 1920 } 1921 1922 SimpleDateFormat dateFormater = null; 1923 if (verbose) { 1924 dateFormater = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss zzz", Locale.US); 1925 } 1926 else { 1927 dateFormater = new SimpleDateFormat("yyyy-MM-dd HH:mm zzz", Locale.US); 1928 } 1929 1930 if (timeZoneId != null) { 1931 dateFormater.setTimeZone(TimeZone.getTimeZone(timeZoneId)); 1932 } 1933 String dateString = dateFormater.format(date); 1934 // Most TimeZones are 3 or 4 characters; GMT offsets (e.g. GMT-07:00) are 9, so lets remove the "GMT" part to make it 6 1935 // to fit better 1936 Matcher m = GMT_OFFSET_SHORTEN_PATTERN.matcher(dateString); 1937 if (m.matches() && m.groupCount() == 2) { 1938 dateString = m.group(1) + m.group(2); 1939 } 1940 return dateString; 1941 } 1942 1943 private void validateCommand(CommandLine commandLine) throws OozieCLIException { 1944 String[] args = commandLine.getArgs(); 1945 if (args.length != 1) { 1946 throw new OozieCLIException("One file must be specified"); 1947 } 1948 try { 1949 XOozieClient wc = createXOozieClient(commandLine); 1950 String result = wc.validateXML(args[0].toString()); 1951 if (result == null) { 1952 // TODO This is only for backward compatibility. Need to remove after 4.2.0 higher version. 1953 System.out.println("Using client-side validation. Check out Oozie server version."); 1954 validateCommandV41(commandLine); 1955 return; 1956 } 1957 System.out.println(result); 1958 } catch (OozieClientException e) { 1959 throw new OozieCLIException(e.getMessage(), e); 1960 } 1961 } 1962 1963 /** 1964 * Validate on client-side. This is only for backward compatibility. Need to removed after <tt>4.2.0</tt> higher version. 1965 * @param commandLine 1966 * @throws OozieCLIException 1967 */ 1968 @Deprecated 1969 @VisibleForTesting 1970 void validateCommandV41(CommandLine commandLine) throws OozieCLIException { 1971 String[] args = commandLine.getArgs(); 1972 if (args.length != 1) { 1973 throw new OozieCLIException("One file must be specified"); 1974 } 1975 File file = new File(args[0]); 1976 if (file.exists()) { 1977 try { 1978 List<StreamSource> sources = new ArrayList<StreamSource>(); 1979 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1980 "oozie-workflow-0.1.xsd"))); 1981 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1982 "shell-action-0.1.xsd"))); 1983 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1984 "shell-action-0.2.xsd"))); 1985 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1986 "shell-action-0.3.xsd"))); 1987 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1988 "email-action-0.1.xsd"))); 1989 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1990 "email-action-0.2.xsd"))); 1991 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1992 "distcp-action-0.1.xsd"))); 1993 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1994 "distcp-action-0.2.xsd"))); 1995 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1996 "oozie-workflow-0.2.xsd"))); 1997 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 1998 "oozie-workflow-0.2.5.xsd"))); 1999 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2000 "oozie-workflow-0.3.xsd"))); 2001 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2002 "oozie-workflow-0.4.xsd"))); 2003 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2004 "oozie-workflow-0.4.5.xsd"))); 2005 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2006 "oozie-workflow-0.5.xsd"))); 2007 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2008 "oozie-coordinator-0.1.xsd"))); 2009 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2010 "oozie-coordinator-0.2.xsd"))); 2011 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2012 "oozie-coordinator-0.3.xsd"))); 2013 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2014 "oozie-coordinator-0.4.xsd"))); 2015 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2016 "oozie-bundle-0.1.xsd"))); 2017 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2018 "oozie-bundle-0.2.xsd"))); 2019 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2020 "oozie-sla-0.1.xsd"))); 2021 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2022 "oozie-sla-0.2.xsd"))); 2023 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2024 "hive-action-0.2.xsd"))); 2025 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2026 "hive-action-0.3.xsd"))); 2027 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2028 "hive-action-0.4.xsd"))); 2029 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2030 "hive-action-0.5.xsd"))); 2031 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2032 "hive-action-0.6.xsd"))); 2033 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2034 "sqoop-action-0.2.xsd"))); 2035 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2036 "sqoop-action-0.3.xsd"))); 2037 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2038 "sqoop-action-0.4.xsd"))); 2039 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2040 "ssh-action-0.1.xsd"))); 2041 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2042 "ssh-action-0.2.xsd"))); 2043 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2044 "hive2-action-0.1.xsd"))); 2045 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2046 "hive2-action-0.2.xsd"))); 2047 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2048 "spark-action-0.1.xsd"))); 2049 sources.add(new StreamSource(Thread.currentThread().getContextClassLoader().getResourceAsStream( 2050 "spark-action-0.2.xsd"))); 2051 SchemaFactory factory = SchemaFactory.newInstance(XMLConstants.W3C_XML_SCHEMA_NS_URI); 2052 Schema schema = factory.newSchema(sources.toArray(new StreamSource[sources.size()])); 2053 Validator validator = schema.newValidator(); 2054 validator.validate(new StreamSource(new FileReader(file))); 2055 System.out.println("Valid workflow-app"); 2056 } 2057 catch (Exception ex) { 2058 throw new OozieCLIException("Invalid app definition, " + ex.toString(), ex); 2059 } 2060 } 2061 else { 2062 throw new OozieCLIException("File does not exists"); 2063 } 2064 } 2065 2066 private void scriptLanguageCommand(CommandLine commandLine, String jobType) throws IOException, OozieCLIException { 2067 List<String> args = commandLine.getArgList(); 2068 if (args.size() > 0) { 2069 // checking if args starts with -X (because CLIParser cannot check this) 2070 if (!args.get(0).equals("-X")) { 2071 throw new OozieCLIException("Unrecognized option: " + args.get(0) + " Expecting -X"); 2072 } 2073 args.remove(0); 2074 } 2075 2076 if (!commandLine.hasOption(SCRIPTFILE_OPTION)) { 2077 throw new OozieCLIException("Need to specify -file <scriptfile>"); 2078 } 2079 2080 if (!commandLine.hasOption(CONFIG_OPTION)) { 2081 throw new OozieCLIException("Need to specify -config <configfile>"); 2082 } 2083 2084 try { 2085 XOozieClient wc = createXOozieClient(commandLine); 2086 Properties conf = getConfiguration(wc, commandLine); 2087 String script = commandLine.getOptionValue(SCRIPTFILE_OPTION); 2088 List<String> paramsList = new ArrayList<String>(); 2089 if (commandLine.hasOption("P")) { 2090 Properties params = commandLine.getOptionProperties("P"); 2091 for (String key : params.stringPropertyNames()) { 2092 paramsList.add(key + "=" + params.getProperty(key)); 2093 } 2094 } 2095 System.out.println(JOB_ID_PREFIX + wc.submitScriptLanguage(conf, script, args.toArray(new String[args.size()]), 2096 paramsList.toArray(new String[paramsList.size()]), jobType)); 2097 } 2098 catch (OozieClientException ex) { 2099 throw new OozieCLIException(ex.toString(), ex); 2100 } 2101 } 2102 2103 private void sqoopCommand(CommandLine commandLine) throws IOException, OozieCLIException { 2104 List<String> args = commandLine.getArgList(); 2105 if (args.size() > 0) { 2106 // checking if args starts with -X (because CLIParser cannot check this) 2107 if (!args.get(0).equals("-X")) { 2108 throw new OozieCLIException("Unrecognized option: " + args.get(0) + " Expecting -X"); 2109 } 2110 args.remove(0); 2111 } 2112 2113 if (!commandLine.hasOption(SQOOP_COMMAND_OPTION)) { 2114 throw new OozieCLIException("Need to specify -command"); 2115 } 2116 2117 if (!commandLine.hasOption(CONFIG_OPTION)) { 2118 throw new OozieCLIException("Need to specify -config <configfile>"); 2119 } 2120 2121 try { 2122 XOozieClient wc = createXOozieClient(commandLine); 2123 Properties conf = getConfiguration(wc, commandLine); 2124 String[] command = commandLine.getOptionValues(SQOOP_COMMAND_OPTION); 2125 System.out.println(JOB_ID_PREFIX + wc.submitSqoop(conf, command, args.toArray(new String[args.size()]))); 2126 } 2127 catch (OozieClientException ex) { 2128 throw new OozieCLIException(ex.toString(), ex); 2129 } 2130 } 2131 2132 private void infoCommand(CommandLine commandLine) throws OozieCLIException { 2133 for (Option option : commandLine.getOptions()) { 2134 String opt = option.getOpt(); 2135 if (opt.equals(INFO_TIME_ZONES_OPTION)) { 2136 printAvailableTimeZones(); 2137 } 2138 } 2139 } 2140 2141 private void printAvailableTimeZones() { 2142 System.out.println("The format is \"SHORT_NAME (ID)\"\nGive the ID to the -timezone argument"); 2143 System.out.println("GMT offsets can also be used (e.g. GMT-07:00, GMT-0700, GMT+05:30, GMT+0530)"); 2144 System.out.println("Available Time Zones:"); 2145 for (String tzId : TimeZone.getAvailableIDs()) { 2146 // skip id's that are like "Etc/GMT+01:00" because their display names are like "GMT-01:00", which is confusing 2147 if (!tzId.startsWith("Etc/GMT")) { 2148 TimeZone tZone = TimeZone.getTimeZone(tzId); 2149 System.out.println(" " + tZone.getDisplayName(false, TimeZone.SHORT) + " (" + tzId + ")"); 2150 } 2151 } 2152 } 2153 2154 2155 private void mrCommand(CommandLine commandLine) throws IOException, OozieCLIException { 2156 try { 2157 XOozieClient wc = createXOozieClient(commandLine); 2158 Properties conf = getConfiguration(wc, commandLine); 2159 2160 String mapper = conf.getProperty(MAPRED_MAPPER, conf.getProperty(MAPRED_MAPPER_2)); 2161 if (mapper == null) { 2162 throw new OozieCLIException("mapper (" + MAPRED_MAPPER + " or " + MAPRED_MAPPER_2 + ") must be specified in conf"); 2163 } 2164 2165 String reducer = conf.getProperty(MAPRED_REDUCER, conf.getProperty(MAPRED_REDUCER_2)); 2166 if (reducer == null) { 2167 throw new OozieCLIException("reducer (" + MAPRED_REDUCER + " or " + MAPRED_REDUCER_2 2168 + ") must be specified in conf"); 2169 } 2170 2171 String inputDir = conf.getProperty(MAPRED_INPUT); 2172 if (inputDir == null) { 2173 throw new OozieCLIException("input dir (" + MAPRED_INPUT +") must be specified in conf"); 2174 } 2175 2176 String outputDir = conf.getProperty(MAPRED_OUTPUT); 2177 if (outputDir == null) { 2178 throw new OozieCLIException("output dir (" + MAPRED_OUTPUT +") must be specified in conf"); 2179 } 2180 2181 System.out.println(JOB_ID_PREFIX + wc.submitMapReduce(conf)); 2182 } 2183 catch (OozieClientException ex) { 2184 throw new OozieCLIException(ex.toString(), ex); 2185 } 2186 } 2187 2188 private String getFirstMissingDependencies(CoordinatorAction action) { 2189 StringBuilder allDeps = new StringBuilder(); 2190 String missingDep = action.getMissingDependencies(); 2191 boolean depExists = false; 2192 if (missingDep != null && !missingDep.isEmpty()) { 2193 allDeps.append(missingDep.split(INSTANCE_SEPARATOR)[0]); 2194 depExists = true; 2195 } 2196 String pushDeps = action.getPushMissingDependencies(); 2197 if (pushDeps != null && !pushDeps.isEmpty()) { 2198 if(depExists) { 2199 allDeps.append(INSTANCE_SEPARATOR); 2200 } 2201 allDeps.append(pushDeps.split(INSTANCE_SEPARATOR)[0]); 2202 } 2203 return allDeps.toString(); 2204 } 2205 2206 private void printMetrics(OozieClient.Metrics metrics) { 2207 System.out.println("COUNTERS"); 2208 System.out.println("--------"); 2209 Map<String, Long> counters = new TreeMap<String, Long>(metrics.getCounters()); 2210 for (Map.Entry<String, Long> ent : counters.entrySet()) { 2211 System.out.println(ent.getKey() + " : " + ent.getValue()); 2212 } 2213 System.out.println("\nGAUGES"); 2214 System.out.println("------"); 2215 Map<String, Object> gauges = new TreeMap<String, Object>(metrics.getGauges()); 2216 for (Map.Entry<String, Object> ent : gauges.entrySet()) { 2217 System.out.println(ent.getKey() + " : " + ent.getValue()); 2218 } 2219 System.out.println("\nTIMERS"); 2220 System.out.println("------"); 2221 Map<String, OozieClient.Metrics.Timer> timers = new TreeMap<String, OozieClient.Metrics.Timer>(metrics.getTimers()); 2222 for (Map.Entry<String, OozieClient.Metrics.Timer> ent : timers.entrySet()) { 2223 System.out.println(ent.getKey()); 2224 System.out.println(ent.getValue()); 2225 } 2226 System.out.println("\nHISTOGRAMS"); 2227 System.out.println("----------"); 2228 Map<String, OozieClient.Metrics.Histogram> histograms = 2229 new TreeMap<String, OozieClient.Metrics.Histogram>(metrics.getHistograms()); 2230 for (Map.Entry<String, OozieClient.Metrics.Histogram> ent : histograms.entrySet()) { 2231 System.out.println(ent.getKey()); 2232 System.out.println(ent.getValue()); 2233 } 2234 } 2235 2236 private void printInstrumentation(OozieClient.Instrumentation instrumentation) { 2237 System.out.println("COUNTERS"); 2238 System.out.println("--------"); 2239 Map<String, Long> counters = new TreeMap<String, Long>(instrumentation.getCounters()); 2240 for (Map.Entry<String, Long> ent : counters.entrySet()) { 2241 System.out.println(ent.getKey() + " : " + ent.getValue()); 2242 } 2243 System.out.println("\nVARIABLES"); 2244 System.out.println("---------"); 2245 Map<String, Object> variables = new TreeMap<String, Object>(instrumentation.getVariables()); 2246 for (Map.Entry<String, Object> ent : variables.entrySet()) { 2247 System.out.println(ent.getKey() + " : " + ent.getValue()); 2248 } 2249 System.out.println("\nSAMPLERS"); 2250 System.out.println("---------"); 2251 Map<String, Double> samplers = new TreeMap<String, Double>(instrumentation.getSamplers()); 2252 for (Map.Entry<String, Double> ent : samplers.entrySet()) { 2253 System.out.println(ent.getKey() + " : " + ent.getValue()); 2254 } 2255 System.out.println("\nTIMERS"); 2256 System.out.println("---------"); 2257 Map<String, OozieClient.Instrumentation.Timer> timers = 2258 new TreeMap<String, OozieClient.Instrumentation.Timer>(instrumentation.getTimers()); 2259 for (Map.Entry<String, OozieClient.Instrumentation.Timer> ent : timers.entrySet()) { 2260 System.out.println(ent.getKey()); 2261 System.out.println(ent.getValue()); 2262 } 2263 } 2264}