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