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.coord; 020 021import com.google.common.collect.Lists; 022import org.apache.commons.lang.StringUtils; 023import org.apache.hadoop.conf.Configuration; 024import org.apache.oozie.ErrorCode; 025import org.apache.oozie.client.OozieClient; 026import org.apache.oozie.command.CommandException; 027import org.apache.oozie.dependency.URIHandler; 028import org.apache.oozie.dependency.URIHandler.Context; 029import org.apache.oozie.service.Services; 030import org.apache.oozie.service.URIHandlerService; 031import org.apache.oozie.util.DateUtils; 032import org.apache.oozie.util.ELEvaluator; 033import org.apache.oozie.util.ParamChecker; 034import org.apache.oozie.util.XLog; 035 036import java.net.URI; 037import java.util.ArrayList; 038import java.util.Calendar; 039import java.util.Date; 040import java.util.GregorianCalendar; 041import java.util.List; 042import java.util.TimeZone; 043 044/** 045 * This class implements the EL function related to coordinator 046 */ 047 048public class CoordELFunctions { 049 final public static String DATASET = "oozie.coord.el.dataset.bean"; 050 final public static String COORD_ACTION = "oozie.coord.el.app.bean"; 051 final public static String CONFIGURATION = "oozie.coord.el.conf"; 052 final public static String LATEST_EL_USE_CURRENT_TIME = "oozie.service.ELService.latest-el.use-current-time"; 053 // INSTANCE_SEPARATOR is used to separate multiple directories into one tag. 054 final public static String INSTANCE_SEPARATOR = "#"; 055 final public static String DIR_SEPARATOR = ","; 056 // TODO: in next release, support flexibility 057 private static String END_OF_OPERATION_INDICATOR_FILE = "_SUCCESS"; 058 059 public static final long MINUTE_MSEC = 60 * 1000L; 060 public static final long HOUR_MSEC = 60 * MINUTE_MSEC; 061 public static final long DAY_MSEC = 24 * HOUR_MSEC; 062 public static final long MONTH_MSEC = 30 * DAY_MSEC; 063 public static final long YEAR_MSEC = 365 * DAY_MSEC; 064 065 /** 066 * Used in defining the frequency in 'day' unit. <p/> domain: <code> val > 0</code> and should be integer. 067 * 068 * @param val frequency in number of days. 069 * @return number of days and also set the frequency timeunit to "day" 070 */ 071 public static int ph1_coord_days(int val) { 072 val = ParamChecker.checkGTZero(val, "n"); 073 ELEvaluator eval = ELEvaluator.getCurrent(); 074 eval.setVariable("timeunit", TimeUnit.DAY); 075 eval.setVariable("endOfDuration", TimeUnit.NONE); 076 return val; 077 } 078 079 /** 080 * Used in defining the frequency in 'month' unit. <p/> domain: <code> val > 0</code> and should be integer. 081 * 082 * @param val frequency in number of months. 083 * @return number of months and also set the frequency timeunit to "month" 084 */ 085 public static int ph1_coord_months(int val) { 086 val = ParamChecker.checkGTZero(val, "n"); 087 ELEvaluator eval = ELEvaluator.getCurrent(); 088 eval.setVariable("timeunit", TimeUnit.MONTH); 089 eval.setVariable("endOfDuration", TimeUnit.NONE); 090 return val; 091 } 092 093 /** 094 * Used in defining the frequency in 'hour' unit. <p/> parameter value domain: <code> val > 0</code> and should 095 * be integer. 096 * 097 * @param val frequency in number of hours. 098 * @return number of minutes and also set the frequency timeunit to "minute" 099 */ 100 public static int ph1_coord_hours(int val) { 101 val = ParamChecker.checkGTZero(val, "n"); 102 ELEvaluator eval = ELEvaluator.getCurrent(); 103 eval.setVariable("timeunit", TimeUnit.MINUTE); 104 eval.setVariable("endOfDuration", TimeUnit.NONE); 105 return val * 60; 106 } 107 108 /** 109 * Used in defining the frequency in 'minute' unit. <p/> domain: <code> val > 0</code> and should be integer. 110 * 111 * @param val frequency in number of minutes. 112 * @return number of minutes and also set the frequency timeunit to "minute" 113 */ 114 public static int ph1_coord_minutes(int val) { 115 val = ParamChecker.checkGTZero(val, "n"); 116 ELEvaluator eval = ELEvaluator.getCurrent(); 117 eval.setVariable("timeunit", TimeUnit.MINUTE); 118 eval.setVariable("endOfDuration", TimeUnit.NONE); 119 return val; 120 } 121 122 /** 123 * Used in defining the frequency in 'day' unit and specify the "end of day" property. <p/> Every instance will 124 * start at 00:00 hour of each day. <p/> domain: <code> val > 0</code> and should be integer. 125 * 126 * @param val frequency in number of days. 127 * @return number of days and also set the frequency timeunit to "day" and end_of_duration flag to "day" 128 */ 129 public static int ph1_coord_endOfDays(int val) { 130 val = ParamChecker.checkGTZero(val, "n"); 131 ELEvaluator eval = ELEvaluator.getCurrent(); 132 eval.setVariable("timeunit", TimeUnit.DAY); 133 eval.setVariable("endOfDuration", TimeUnit.END_OF_DAY); 134 return val; 135 } 136 137 /** 138 * Used in defining the frequency in 'month' unit and specify the "end of month" property. <p/> Every instance will 139 * start at first day of each month at 00:00 hour. <p/> domain: <code> val > 0</code> and should be integer. 140 * 141 * @param val: frequency in number of months. 142 * @return number of months and also set the frequency timeunit to "month" and end_of_duration flag to "month" 143 */ 144 public static int ph1_coord_endOfMonths(int val) { 145 val = ParamChecker.checkGTZero(val, "n"); 146 ELEvaluator eval = ELEvaluator.getCurrent(); 147 eval.setVariable("timeunit", TimeUnit.MONTH); 148 eval.setVariable("endOfDuration", TimeUnit.END_OF_MONTH); 149 return val; 150 } 151 152 /** 153 * Calculate the difference of timezone offset in minutes between dataset and coordinator job. <p/> Depends on: <p/> 154 * 1. Timezone of both dataset and job <p/> 2. Action creation Time 155 * 156 * @return difference in minutes (DataSet TZ Offset - Application TZ offset) 157 */ 158 public static int ph2_coord_tzOffset() { 159 long actionCreationTime = getActionCreationtime().getTime(); 160 TimeZone dsTZ = ParamChecker.notNull(getDatasetTZ(), "DatasetTZ"); 161 TimeZone jobTZ = ParamChecker.notNull(getJobTZ(), "JobTZ"); 162 return (dsTZ.getOffset(actionCreationTime) - jobTZ.getOffset(actionCreationTime)) / (1000 * 60); 163 } 164 165 public static int ph3_coord_tzOffset() { 166 return ph2_coord_tzOffset(); 167 } 168 169 /** 170 * Returns a date string that is offset from 'strBaseDate' by the amount specified. The unit can be one of 171 * DAY, MONTH, HOUR, MINUTE, MONTH. 172 * 173 * @param strBaseDate The base date 174 * @param offset any number 175 * @param unit one of DAY, MONTH, HOUR, MINUTE, MONTH 176 * @return the offset date string 177 * @throws Exception 178 */ 179 public static String ph2_coord_dateOffset(String strBaseDate, int offset, String unit) throws Exception { 180 Calendar baseCalDate = DateUtils.getCalendar(strBaseDate); 181 StringBuilder buffer = new StringBuilder(); 182 baseCalDate.add(TimeUnit.valueOf(unit).getCalendarUnit(), offset); 183 buffer.append(DateUtils.formatDateOozieTZ(baseCalDate)); 184 return buffer.toString(); 185 } 186 187 public static String ph3_coord_dateOffset(String strBaseDate, int offset, String unit) throws Exception { 188 return ph2_coord_dateOffset(strBaseDate, offset, unit); 189 } 190 191 /** 192 * Returns a date string that is offset from 'strBaseDate' by the difference from Oozie processing timezone to the given 193 * timezone. It will account for daylight saving time based on the given 'strBaseDate' and 'timezone'. 194 * 195 * @param strBaseDate The base date 196 * @param timezone 197 * @return the offset date string 198 * @throws Exception 199 */ 200 public static String ph2_coord_dateTzOffset(String strBaseDate, String timezone) throws Exception { 201 Calendar baseCalDate = DateUtils.getCalendar(strBaseDate); 202 StringBuilder buffer = new StringBuilder(); 203 baseCalDate.setTimeZone(DateUtils.getTimeZone(timezone)); 204 buffer.append(DateUtils.formatDate(baseCalDate)); 205 return buffer.toString(); 206 } 207 208 public static String ph3_coord_dateTzOffset(String strBaseDate, String timezone) throws Exception{ 209 return ph2_coord_dateTzOffset(strBaseDate, timezone); 210 } 211 212 /** 213 * Determine the date-time in Oozie processing timezone of n-th future available dataset instance 214 * from nominal Time but not beyond the instance specified as 'instance. 215 * <p/> 216 * It depends on: 217 * <p/> 218 * 1. Data set frequency 219 * <p/> 220 * 2. Data set Time unit (day, month, minute) 221 * <p/> 222 * 3. Data set Time zone/DST 223 * <p/> 224 * 4. End Day/Month flag 225 * <p/> 226 * 5. Data set initial instance 227 * <p/> 228 * 6. Action Creation Time 229 * <p/> 230 * 7. Existence of dataset's directory 231 * 232 * @param n :instance count 233 * <p/> 234 * domain: n >= 0, n is integer 235 * @param instance: How many future instance it should check? value should 236 * be >=0 237 * @return date-time in Oozie processing timezone of the n-th instance 238 * <p/> 239 * @throws Exception 240 */ 241 public static String ph3_coord_future(int n, int instance) throws Exception { 242 ParamChecker.checkGEZero(n, "future:n"); 243 ParamChecker.checkGTZero(instance, "future:instance"); 244 if (isSyncDataSet()) {// For Sync Dataset 245 return coord_future_sync(n, instance); 246 } 247 else { 248 throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet"); 249 } 250 } 251 252 /** 253 * Determine the date-time in Oozie processing timezone of the future available dataset instances 254 * from start to end offsets from nominal Time but not beyond the instance specified as 'instance'. 255 * <p/> 256 * It depends on: 257 * <p/> 258 * 1. Data set frequency 259 * <p/> 260 * 2. Data set Time unit (day, month, minute) 261 * <p/> 262 * 3. Data set Time zone/DST 263 * <p/> 264 * 4. End Day/Month flag 265 * <p/> 266 * 5. Data set initial instance 267 * <p/> 268 * 6. Action Creation Time 269 * <p/> 270 * 7. Existence of dataset's directory 271 * 272 * @param start : start instance offset 273 * <p/> 274 * domain: start >= 0, start is integer 275 * @param end : end instance offset 276 * <p/> 277 * domain: end >= 0, end is integer 278 * @param instance: How many future instance it should check? value should 279 * be >=0 280 * @return date-time in Oozie processing timezone of the instances from start to end offsets 281 * delimited by comma. 282 * <p/> 283 * @throws Exception 284 */ 285 public static String ph3_coord_futureRange(int start, int end, int instance) throws Exception { 286 ParamChecker.checkGEZero(start, "future:n"); 287 ParamChecker.checkGEZero(end, "future:n"); 288 ParamChecker.checkGTZero(instance, "future:instance"); 289 if (isSyncDataSet()) {// For Sync Dataset 290 return coord_futureRange_sync(start, end, instance); 291 } 292 else { 293 throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet"); 294 } 295 } 296 297 private static String coord_future_sync(int n, int instance) throws Exception { 298 return coord_futureRange_sync(n, n, instance); 299 } 300 301 private static String coord_futureRange_sync(int startOffset, int endOffset, int instance) throws Exception { 302 final XLog LOG = XLog.getLog(CoordELFunctions.class); 303 final Thread currentThread = Thread.currentThread(); 304 ELEvaluator eval = ELEvaluator.getCurrent(); 305 String retVal = ""; 306 int datasetFrequency = (int) getDSFrequency();// in minutes 307 TimeUnit dsTimeUnit = getDSTimeUnit(); 308 int[] instCount = new int[1]; 309 Calendar nominalInstanceCal = getCurrentInstance(getActionCreationtime(), instCount); 310 StringBuilder resolvedInstances = new StringBuilder(); 311 StringBuilder resolvedURIPaths = new StringBuilder(); 312 if (nominalInstanceCal != null) { 313 Calendar initInstance = getInitialInstanceCal(); 314 nominalInstanceCal = (Calendar) initInstance.clone(); 315 nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), instCount[0] * datasetFrequency); 316 317 SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET); 318 if (ds == null) { 319 throw new RuntimeException("Associated Dataset should be defined with key " + DATASET); 320 } 321 String uriTemplate = ds.getUriTemplate(); 322 Configuration conf = (Configuration) eval.getVariable(CONFIGURATION); 323 if (conf == null) { 324 throw new RuntimeException("Associated Configuration should be defined with key " + CONFIGURATION); 325 } 326 int available = 0, checkedInstance = 0; 327 boolean resolved = false; 328 String user = ParamChecker 329 .notEmpty((String) eval.getVariable(OozieClient.USER_NAME), OozieClient.USER_NAME); 330 String doneFlag = ds.getDoneFlag(); 331 URIHandlerService uriService = Services.get().get(URIHandlerService.class); 332 URIHandler uriHandler = null; 333 Context uriContext = null; 334 try { 335 while (instance >= checkedInstance && !currentThread.isInterrupted()) { 336 ELEvaluator uriEval = getUriEvaluator(nominalInstanceCal); 337 String uriPath = uriEval.evaluate(uriTemplate, String.class); 338 if (uriHandler == null) { 339 URI uri = new URI(uriPath); 340 uriHandler = uriService.getURIHandler(uri); 341 uriContext = uriHandler.getContext(uri, conf, user); 342 } 343 String uriWithDoneFlag = uriHandler.getURIWithDoneFlag(uriPath, doneFlag); 344 if (uriHandler.exists(new URI(uriWithDoneFlag), uriContext)) { 345 if (available == endOffset) { 346 LOG.debug("Matched future(" + available + "): " + uriWithDoneFlag); 347 resolved = true; 348 resolvedInstances.append(DateUtils.formatDateOozieTZ(nominalInstanceCal)); 349 resolvedURIPaths.append(uriPath); 350 retVal = resolvedInstances.toString(); 351 eval.setVariable("resolved_path", resolvedURIPaths.toString()); 352 break; 353 } 354 else if (available >= startOffset) { 355 LOG.debug("Matched future(" + available + "): " + uriWithDoneFlag); 356 resolvedInstances.append(DateUtils.formatDateOozieTZ(nominalInstanceCal)).append( 357 INSTANCE_SEPARATOR); 358 resolvedURIPaths.append(uriPath).append(INSTANCE_SEPARATOR); 359 } 360 available++; 361 } 362 // nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), datasetFrequency); 363 nominalInstanceCal = (Calendar) initInstance.clone(); 364 instCount[0]++; 365 nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), instCount[0] * datasetFrequency); 366 checkedInstance++; 367 // DateUtils.moveToEnd(nominalInstanceCal, getDSEndOfFlag()); 368 } 369 } 370 finally { 371 if (uriContext != null) { 372 uriContext.destroy(); 373 } 374 } 375 if (!resolved) { 376 // return unchanged future function with variable 'is_resolved' 377 // to 'false' 378 eval.setVariable("is_resolved", Boolean.FALSE); 379 if (startOffset == endOffset) { 380 retVal = "${coord:future(" + startOffset + ", " + instance + ")}"; 381 } 382 else { 383 retVal = "${coord:futureRange(" + startOffset + ", " + endOffset + ", " + instance + ")}"; 384 } 385 } 386 else { 387 eval.setVariable("is_resolved", Boolean.TRUE); 388 } 389 } 390 else {// No feasible nominal time 391 eval.setVariable("is_resolved", Boolean.TRUE); 392 retVal = ""; 393 } 394 return retVal; 395 } 396 397 /** 398 * Return nominal time or Action Creation Time. 399 * <p/> 400 * 401 * @return coordinator action creation or materialization date time 402 * @throws Exception if unable to format the Date object to String 403 */ 404 public static String ph2_coord_nominalTime() throws Exception { 405 ELEvaluator eval = ELEvaluator.getCurrent(); 406 SyncCoordAction action = ParamChecker.notNull((SyncCoordAction) eval.getVariable(COORD_ACTION), 407 "Coordinator Action"); 408 return DateUtils.formatDateOozieTZ(action.getNominalTime()); 409 } 410 411 public static String ph3_coord_nominalTime() throws Exception { 412 return ph2_coord_nominalTime(); 413 } 414 415 /** 416 * Convert from standard date-time formatting to a desired format. 417 * <p/> 418 * @param dateTimeStr - A timestamp in standard (ISO8601) format. 419 * @param format - A string representing the desired format. 420 * @return coordinator action creation or materialization date time 421 * @throws Exception if unable to format the Date object to String 422 */ 423 public static String ph2_coord_formatTime(String dateTimeStr, String format) 424 throws Exception { 425 Date dateTime = DateUtils.parseDateOozieTZ(dateTimeStr); 426 return DateUtils.formatDateCustom(dateTime, format); 427 } 428 429 public static String ph3_coord_formatTime(String dateTimeStr, String format) 430 throws Exception { 431 return ph2_coord_formatTime(dateTimeStr, format); 432 } 433 434 /** 435 * Return Action Id. <p/> 436 * Convert from standard date-time formatting to a Unix epoch time. 437 * <p/> 438 * @param dateTimeStr - A timestamp in standard (ISO8601) format. 439 * @param millis - "true" to include millis; otherwise will only include seconds 440 * @return coordinator action creation or materialization date time 441 * @throws Exception if unable to format the Date object to String 442 */ 443 public static String ph2_coord_epochTime(String dateTimeStr, String millis) 444 throws Exception { 445 Date dateTime = DateUtils.parseDateOozieTZ(dateTimeStr); 446 return DateUtils.formatDateEpoch(dateTime, Boolean.valueOf(millis)); 447 } 448 449 public static String ph3_coord_epochTime(String dateTimeStr, String millis) 450 throws Exception { 451 return ph2_coord_epochTime(dateTimeStr, millis); 452 } 453 454 /** 455 * Return Action Id. <p> 456 * 457 * @return coordinator action Id 458 */ 459 public static String ph2_coord_actionId() throws Exception { 460 ELEvaluator eval = ELEvaluator.getCurrent(); 461 SyncCoordAction action = ParamChecker.notNull((SyncCoordAction) eval.getVariable(COORD_ACTION), 462 "Coordinator Action"); 463 return action.getActionId(); 464 } 465 466 public static String ph3_coord_actionId() throws Exception { 467 return ph2_coord_actionId(); 468 } 469 470 /** 471 * Return Job Name. <p/> 472 * 473 * @return coordinator name 474 */ 475 public static String ph2_coord_name() throws Exception { 476 ELEvaluator eval = ELEvaluator.getCurrent(); 477 SyncCoordAction action = ParamChecker.notNull((SyncCoordAction) eval.getVariable(COORD_ACTION), 478 "Coordinator Action"); 479 return action.getName(); 480 } 481 482 public static String ph3_coord_name() throws Exception { 483 return ph2_coord_name(); 484 } 485 486 /** 487 * Return Action Start time. <p/> 488 * 489 * @return coordinator action start time 490 * @throws Exception if unable to format the Date object to String 491 */ 492 public static String ph2_coord_actualTime() throws Exception { 493 ELEvaluator eval = ELEvaluator.getCurrent(); 494 SyncCoordAction coordAction = (SyncCoordAction) eval.getVariable(COORD_ACTION); 495 if (coordAction == null) { 496 throw new RuntimeException("Associated Application instance should be defined with key " + COORD_ACTION); 497 } 498 return DateUtils.formatDateOozieTZ(coordAction.getActualTime()); 499 } 500 501 public static String ph3_coord_actualTime() throws Exception { 502 return ph2_coord_actualTime(); 503 } 504 505 /** 506 * Used to specify a list of URI's that are used as input dir to the workflow job. <p/> Look for two evaluator-level 507 * variables <p/> A) .datain.<DATAIN_NAME> B) .datain.<DATAIN_NAME>.unresolved <p/> A defines the current list of 508 * URI. <p/> B defines whether there are any unresolved EL-function (i.e latest) <p/> If there are something 509 * unresolved, this function will echo back the original function <p/> otherwise it sends the uris. 510 * 511 * @param dataInName : Datain name 512 * @return the list of URI's separated by INSTANCE_SEPARATOR <p/> if there are unresolved EL function (i.e. latest) 513 * , echo back <p/> the function without resolving the function. 514 */ 515 public static String ph3_coord_dataIn(String dataInName) { 516 String uris = ""; 517 ELEvaluator eval = ELEvaluator.getCurrent(); 518 uris = (String) eval.getVariable(".datain." + dataInName); 519 Boolean unresolved = (Boolean) eval.getVariable(".datain." + dataInName + ".unresolved"); 520 if (unresolved != null && unresolved.booleanValue() == true) { 521 return "${coord:dataIn('" + dataInName + "')}"; 522 } 523 return uris; 524 } 525 526 /** 527 * Used to specify a list of URI's that are output dir of the workflow job. <p/> Look for one evaluator-level 528 * variable <p/> dataout.<DATAOUT_NAME> <p/> It defines the current list of URI. <p/> otherwise it sends the uris. 529 * 530 * @param dataOutName : Dataout name 531 * @return the list of URI's separated by INSTANCE_SEPARATOR 532 */ 533 public static String ph3_coord_dataOut(String dataOutName) { 534 String uris = ""; 535 ELEvaluator eval = ELEvaluator.getCurrent(); 536 uris = (String) eval.getVariable(".dataout." + dataOutName); 537 return uris; 538 } 539 540 /** 541 * Determine the date-time in Oozie processing timezone of n-th dataset instance. <p/> It depends on: <p/> 1. 542 * Data set frequency <p/> 2. 543 * Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST <p/> 4. End Day/Month flag <p/> 5. Data 544 * set initial instance <p/> 6. Action Creation Time 545 * 546 * @param n instance count domain: n is integer 547 * @return date-time in Oozie processing timezone of the n-th instance returns 'null' means n-th instance is 548 * earlier than Initial-Instance of DS 549 * @throws Exception 550 */ 551 public static String ph2_coord_current(int n) throws Exception { 552 if (isSyncDataSet()) { // For Sync Dataset 553 return coord_current_sync(n); 554 } 555 else { 556 throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet"); 557 } 558 } 559 560 /** 561 * Determine the date-time in Oozie processing timezone of current dataset instances 562 * from start to end offsets from the nominal time. <p/> It depends 563 * on: <p/> 1. Data set frequency <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST 564 * <p/> 4. End Day/Month flag <p/> 5. Data set initial instance <p/> 6. Action Creation Time 565 * 566 * @param start :start instance offset <p/> domain: start <= 0, start is integer 567 * @param end :end instance offset <p/> domain: end <= 0, end is integer 568 * @return date-time in Oozie processing timezone of the instances from start to end offsets 569 * delimited by comma. <p/> If the current instance time of the dataset based on the Action Creation Time 570 * is earlier than the Initial-Instance of DS an empty string is returned. 571 * If an instance within the range is earlier than Initial-Instance of DS that instance is ignored 572 * @throws Exception 573 */ 574 public static String ph2_coord_currentRange(int start, int end) throws Exception { 575 if (isSyncDataSet()) { // For Sync Dataset 576 return coord_currentRange_sync(start, end); 577 } 578 else { 579 throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet"); 580 } 581 } 582 /** 583 * Determine the date-time in Oozie processing timezone of the given offset from the dataset effective nominal time. <p/> It 584 * depends on: <p> 1. Data set frequency <p/> 2. Data set Time Unit <p/> 3. Data set Time zone/DST 585 * <p/> 4. Data set initial instance <p/> 5. Action Creation Time 586 * 587 * @param n offset amount (integer) 588 * @param timeUnit TimeUnit for offset n ("MINUTE", "HOUR", "DAY", "MONTH", "YEAR") 589 * @return date-time in Oozie processing timezone of the given offset from the dataset effective nominal time 590 * @throws Exception if there was a problem formatting 591 */ 592 public static String ph2_coord_offset(int n, String timeUnit) throws Exception { 593 if (isSyncDataSet()) { // For Sync Dataset 594 return coord_offset_sync(n, timeUnit); 595 } 596 else { 597 throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet"); 598 } 599 } 600 601 /** 602 * Determine how many hours is on the date of n-th dataset instance. <p/> It depends on: <p/> 1. Data set frequency 603 * <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST <p/> 4. End Day/Month flag <p/> 5. 604 * Data set initial instance <p/> 6. Action Creation Time 605 * 606 * @param n instance count <p/> domain: n is integer 607 * @return number of hours on that day <p/> returns -1 means n-th instance is earlier than Initial-Instance of DS 608 * @throws Exception 609 */ 610 public static int ph2_coord_hoursInDay(int n) throws Exception { 611 int datasetFrequency = (int) getDSFrequency(); 612 // /Calendar nominalInstanceCal = 613 // getCurrentInstance(getActionCreationtime()); 614 Calendar nominalInstanceCal = getEffectiveNominalTime(); 615 if (nominalInstanceCal == null) { 616 return -1; 617 } 618 nominalInstanceCal.add(getDSTimeUnit().getCalendarUnit(), datasetFrequency * n); 619 /* 620 * if (nominalInstanceCal.getTime().compareTo(getInitialInstance()) < 0) 621 * { return -1; } 622 */ 623 nominalInstanceCal.setTimeZone(getDatasetTZ());// Use Dataset TZ 624 // DateUtils.moveToEnd(nominalInstanceCal, getDSEndOfFlag()); 625 return DateUtils.hoursInDay(nominalInstanceCal); 626 } 627 628 public static int ph3_coord_hoursInDay(int n) throws Exception { 629 return ph2_coord_hoursInDay(n); 630 } 631 632 /** 633 * Calculate number of days in one month for n-th dataset instance. <p/> It depends on: <p/> 1. Data set frequency . 634 * <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST <p/> 4. End Day/Month flag <p/> 5. 635 * Data set initial instance <p/> 6. Action Creation Time 636 * 637 * @param n instance count. domain: n is integer 638 * @return number of days in that month <p/> returns -1 means n-th instance is earlier than Initial-Instance of DS 639 * @throws Exception 640 */ 641 public static int ph2_coord_daysInMonth(int n) throws Exception { 642 int datasetFrequency = (int) getDSFrequency();// in minutes 643 // Calendar nominalInstanceCal = 644 // getCurrentInstance(getActionCreationtime()); 645 Calendar nominalInstanceCal = getEffectiveNominalTime(); 646 if (nominalInstanceCal == null) { 647 return -1; 648 } 649 nominalInstanceCal.add(getDSTimeUnit().getCalendarUnit(), datasetFrequency * n); 650 /* 651 * if (nominalInstanceCal.getTime().compareTo(getInitialInstance()) < 0) 652 * { return -1; } 653 */ 654 nominalInstanceCal.setTimeZone(getDatasetTZ());// Use Dataset TZ 655 // DateUtils.moveToEnd(nominalInstanceCal, getDSEndOfFlag()); 656 return nominalInstanceCal.getActualMaximum(Calendar.DAY_OF_MONTH); 657 } 658 659 public static int ph3_coord_daysInMonth(int n) throws Exception { 660 return ph2_coord_daysInMonth(n); 661 } 662 663 /** 664 * Determine the date-time in Oozie processing timezone of n-th latest available dataset instance. <p/> It depends 665 * on: <p/> 1. Data set frequency <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST 666 * <p/> 4. End Day/Month flag <p/> 5. Data set initial instance <p/> 6. Action Creation Time <p/> 7. Existence of 667 * dataset's directory 668 * 669 * @param n :instance count <p/> domain: n <= 0, n is integer 670 * @return date-time in Oozie processing timezone of the n-th instance <p/> returns 'null' means n-th instance is 671 * earlier than Initial-Instance of DS 672 * @throws Exception 673 */ 674 public static String ph3_coord_latest(int n) throws Exception { 675 ParamChecker.checkLEZero(n, "latest:n"); 676 if (isSyncDataSet()) {// For Sync Dataset 677 return coord_latest_sync(n); 678 } 679 else { 680 throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet"); 681 } 682 } 683 684 /** 685 * Determine the date-time in Oozie processing timezone of latest available dataset instances 686 * from start to end offsets from the nominal time. <p/> It depends 687 * on: <p/> 1. Data set frequency <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST 688 * <p/> 4. End Day/Month flag <p/> 5. Data set initial instance <p/> 6. Action Creation Time <p/> 7. Existence of 689 * dataset's directory 690 * 691 * @param start :start instance offset <p/> domain: start <= 0, start is integer 692 * @param end :end instance offset <p/> domain: end <= 0, end is integer 693 * @return date-time in Oozie processing timezone of the instances from start to end offsets 694 * delimited by comma. <p/> returns 'null' means start offset instance is 695 * earlier than Initial-Instance of DS 696 * @throws Exception 697 */ 698 public static String ph3_coord_latestRange(int start, int end) throws Exception { 699 ParamChecker.checkLEZero(start, "latest:n"); 700 ParamChecker.checkLEZero(end, "latest:n"); 701 if (isSyncDataSet()) {// For Sync Dataset 702 return coord_latestRange_sync(start, end); 703 } 704 else { 705 throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet"); 706 } 707 } 708 709 /** 710 * Configure an evaluator with data set and application specific information. <p/> Helper method of associating 711 * dataset and application object 712 * 713 * @param evaluator : to set variables 714 * @param ds : Data Set object 715 * @param coordAction : Application instance 716 */ 717 public static void configureEvaluator(ELEvaluator evaluator, SyncCoordDataset ds, SyncCoordAction coordAction) { 718 evaluator.setVariable(COORD_ACTION, coordAction); 719 evaluator.setVariable(DATASET, ds); 720 } 721 722 /** 723 * Helper method to wrap around with "${..}". <p/> 724 * 725 * 726 * @param eval :EL evaluator 727 * @param expr : expression to evaluate 728 * @return Resolved expression or echo back the same expression 729 * @throws Exception 730 */ 731 public static String evalAndWrap(ELEvaluator eval, String expr) throws Exception { 732 try { 733 eval.setVariable(".wrap", null); 734 String result = eval.evaluate(expr, String.class); 735 if (eval.getVariable(".wrap") != null) { 736 return "${" + result + "}"; 737 } 738 else { 739 return result; 740 } 741 } 742 catch (Exception e) { 743 throw new Exception("Unable to evaluate :" + expr + ":\n", e); 744 } 745 } 746 747 // Set of echo functions 748 749 public static String ph1_coord_current_echo(String n) { 750 return echoUnResolved("current", n); 751 } 752 753 public static String ph1_coord_absolute_echo(String date) { 754 return echoUnResolved("absolute", date); 755 } 756 757 public static String ph1_coord_currentRange_echo(String start, String end) { 758 return echoUnResolved("currentRange", start + ", " + end); 759 } 760 761 public static String ph1_coord_offset_echo(String n, String timeUnit) { 762 return echoUnResolved("offset", n + " , " + timeUnit); 763 } 764 765 public static String ph2_coord_current_echo(String n) { 766 return echoUnResolved("current", n); 767 } 768 769 public static String ph2_coord_currentRange_echo(String start, String end) { 770 return echoUnResolved("currentRange", start + ", " + end); 771 } 772 773 public static String ph2_coord_offset_echo(String n, String timeUnit) { 774 return echoUnResolved("offset", n + " , " + timeUnit); 775 } 776 777 public static String ph2_coord_absolute_echo(String date) { 778 return echoUnResolved("absolute", date); 779 } 780 781 public static String ph2_coord_absolute_range(String startInstance, int end) throws Exception { 782 int[] instanceCount = new int[1]; 783 784 // getCurrentInstance() returns null, which means startInstance is less 785 // than initial instance 786 if (getCurrentInstance(DateUtils.getCalendar(startInstance).getTime(), instanceCount) == null) { 787 throw new CommandException(ErrorCode.E1010, 788 "intial-instance should be equal or earlier than the start-instance. intial-instance is " 789 + getInitialInstance() + " and start-instance is " + startInstance); 790 } 791 int[] nominalCount = new int[1]; 792 if (getCurrentInstance(getActionCreationtime(), nominalCount) == null) { 793 throw new CommandException(ErrorCode.E1010, 794 "intial-instance should be equal or earlier than the nominal time. intial-instance is " 795 + getInitialInstance() + " and nominal time is " + getActionCreationtime()); 796 } 797 // getCurrentInstance return offset relative to initial instance. 798 // start instance offset - nominal offset = start offset relative to 799 // nominal time-stamp. 800 int start = instanceCount[0] - nominalCount[0]; 801 if (start > end) { 802 throw new CommandException(ErrorCode.E1010, 803 "start-instance should be equal or earlier than the end-instance. startInstance is " 804 + startInstance + " which is equivalent to current (" + instanceCount[0] 805 + ") but end is specified as current (" + end + ")"); 806 } 807 return ph2_coord_currentRange(start, end); 808 } 809 810 public static String ph1_coord_dateOffset_echo(String n, String offset, String unit) { 811 return echoUnResolved("dateOffset", n + " , " + offset + " , " + unit); 812 } 813 814 public static String ph1_coord_dateTzOffset_echo(String n, String timezone) { 815 return echoUnResolved("dateTzOffset", n + " , " + timezone); 816 } 817 818 public static String ph1_coord_epochTime_echo(String dateTime, String millis) { 819 // Quote the dateTime value since it would contain a ':'. 820 return echoUnResolved("epochTime", "'"+dateTime+"'" + " , " + millis); 821 } 822 823 public static String ph1_coord_formatTime_echo(String dateTime, String format) { 824 // Quote the dateTime value since it would contain a ':'. 825 return echoUnResolved("formatTime", "'"+dateTime+"'" + " , " + format); 826 } 827 828 public static String ph1_coord_latest_echo(String n) { 829 return echoUnResolved("latest", n); 830 } 831 832 public static String ph2_coord_latest_echo(String n) { 833 return ph1_coord_latest_echo(n); 834 } 835 836 public static String ph1_coord_future_echo(String n, String instance) { 837 return echoUnResolved("future", n + ", " + instance + ""); 838 } 839 840 public static String ph2_coord_future_echo(String n, String instance) { 841 return ph1_coord_future_echo(n, instance); 842 } 843 844 public static String ph1_coord_latestRange_echo(String start, String end) { 845 return echoUnResolved("latestRange", start + ", " + end); 846 } 847 848 public static String ph2_coord_latestRange_echo(String start, String end) { 849 return ph1_coord_latestRange_echo(start, end); 850 } 851 852 public static String ph1_coord_futureRange_echo(String start, String end, String instance) { 853 return echoUnResolved("futureRange", start + ", " + end + ", " + instance); 854 } 855 856 public static String ph2_coord_futureRange_echo(String start, String end, String instance) { 857 return ph1_coord_futureRange_echo(start, end, instance); 858 } 859 860 public static String ph1_coord_dataIn_echo(String n) { 861 ELEvaluator eval = ELEvaluator.getCurrent(); 862 String val = (String) eval.getVariable("oozie.dataname." + n); 863 if (val == null || val.equals("data-in") == false) { 864 XLog.getLog(CoordELFunctions.class).error("data_in_name " + n + " is not valid"); 865 throw new RuntimeException("data_in_name " + n + " is not valid"); 866 } 867 return echoUnResolved("dataIn", "'" + n + "'"); 868 } 869 870 public static String ph1_coord_dataOut_echo(String n) { 871 ELEvaluator eval = ELEvaluator.getCurrent(); 872 String val = (String) eval.getVariable("oozie.dataname." + n); 873 if (val == null || val.equals("data-out") == false) { 874 XLog.getLog(CoordELFunctions.class).error("data_out_name " + n + " is not valid"); 875 throw new RuntimeException("data_out_name " + n + " is not valid"); 876 } 877 return echoUnResolved("dataOut", "'" + n + "'"); 878 } 879 880 public static String ph1_coord_nominalTime_echo() { 881 return echoUnResolved("nominalTime", ""); 882 } 883 884 public static String ph1_coord_nominalTime_echo_wrap() { 885 // return "${coord:nominalTime()}"; // no resolution 886 return echoUnResolved("nominalTime", ""); 887 } 888 889 public static String ph1_coord_nominalTime_echo_fixed() { 890 return "2009-03-06T010:00"; // Dummy resolution 891 } 892 893 public static String ph1_coord_actualTime_echo_wrap() { 894 // return "${coord:actualTime()}"; // no resolution 895 return echoUnResolved("actualTime", ""); 896 } 897 898 public static String ph1_coord_actionId_echo() { 899 return echoUnResolved("actionId", ""); 900 } 901 902 public static String ph1_coord_name_echo() { 903 return echoUnResolved("name", ""); 904 } 905 906 // The following echo functions are not used in any phases yet 907 // They are here for future purpose. 908 public static String coord_minutes_echo(String n) { 909 return echoUnResolved("minutes", n); 910 } 911 912 public static String coord_hours_echo(String n) { 913 return echoUnResolved("hours", n); 914 } 915 916 public static String coord_days_echo(String n) { 917 return echoUnResolved("days", n); 918 } 919 920 public static String coord_endOfDay_echo(String n) { 921 return echoUnResolved("endOfDay", n); 922 } 923 924 public static String coord_months_echo(String n) { 925 return echoUnResolved("months", n); 926 } 927 928 public static String coord_endOfMonth_echo(String n) { 929 return echoUnResolved("endOfMonth", n); 930 } 931 932 public static String coord_actualTime_echo() { 933 return echoUnResolved("actualTime", ""); 934 } 935 936 // This echo function will always return "24" for validation only. 937 // This evaluation ****should not**** replace the original XML 938 // Create a temporary string and validate the function 939 // This is **required** for evaluating an expression like 940 // coord:HoursInDay(0) + 3 941 // actual evaluation will happen in phase 2 or phase 3. 942 public static String ph1_coord_hoursInDay_echo(String n) { 943 return "24"; 944 // return echoUnResolved("hoursInDay", n); 945 } 946 947 // This echo function will always return "30" for validation only. 948 // This evaluation ****should not**** replace the original XML 949 // Create a temporary string and validate the function 950 // This is **required** for evaluating an expression like 951 // coord:daysInMonth(0) + 3 952 // actual evaluation will happen in phase 2 or phase 3. 953 public static String ph1_coord_daysInMonth_echo(String n) { 954 // return echoUnResolved("daysInMonth", n); 955 return "30"; 956 } 957 958 // This echo function will always return "3" for validation only. 959 // This evaluation ****should not**** replace the original XML 960 // Create a temporary string and validate the function 961 // This is **required** for evaluating an expression like coord:tzOffset + 2 962 // actual evaluation will happen in phase 2 or phase 3. 963 public static String ph1_coord_tzOffset_echo() { 964 // return echoUnResolved("tzOffset", ""); 965 return "3"; 966 } 967 968 // Local methods 969 /** 970 * @param n 971 * @return n-th instance Date-Time from current instance for data-set <p/> return empty string ("") if the 972 * Action_Creation_time or the n-th instance <p/> is earlier than the Initial_Instance of dataset. 973 * @throws Exception 974 */ 975 private static String coord_current_sync(int n) throws Exception { 976 return coord_currentRange_sync(n, n); 977 } 978 979 private static String coord_currentRange_sync(int start, int end) throws Exception { 980 final XLog LOG = XLog.getLog(CoordELFunctions.class); 981 int datasetFrequency = getDSFrequency();// in minutes 982 TimeUnit dsTimeUnit = getDSTimeUnit(); 983 int[] instCount = new int[1];// used as pass by ref 984 Calendar nominalInstanceCal = getCurrentInstance(getActionCreationtime(), instCount); 985 if (nominalInstanceCal == null) { 986 LOG.warn("If the initial instance of the dataset is later than the nominal time, an empty string is" 987 + " returned. This means that no data is available at the current-instance specified by the user" 988 + " and the user could try modifying his initial-instance to an earlier time."); 989 return ""; 990 } else { 991 Calendar initInstance = getInitialInstanceCal(); 992 // Add in the reverse order - newest instance first. 993 nominalInstanceCal = (Calendar) initInstance.clone(); 994 nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), (instCount[0] + start) * datasetFrequency); 995 List<String> instances = new ArrayList<String>(); 996 for (int i = start; i <= end; i++) { 997 if (nominalInstanceCal.compareTo(initInstance) < 0) { 998 LOG.warn("If the initial instance of the dataset is later than the current-instance specified," 999 + " such as coord:current({0}) in this case, an empty string is returned. This means that" 1000 + " no data is available at the current-instance specified by the user and the user could" 1001 + " try modifying his initial-instance to an earlier time.", start); 1002 } 1003 else { 1004 instances.add(DateUtils.formatDateOozieTZ(nominalInstanceCal)); 1005 } 1006 nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), datasetFrequency); 1007 } 1008 instances = Lists.reverse(instances); 1009 return StringUtils.join(instances, CoordELFunctions.INSTANCE_SEPARATOR); 1010 } 1011 } 1012 1013 /** 1014 * 1015 * @param n offset amount (integer) 1016 * @param timeUnit TimeUnit for offset n ("MINUTE", "HOUR", "DAY", "MONTH", "YEAR") 1017 * @return the offset time from the effective nominal time <p/> return empty string ("") if the Action_Creation_time or the 1018 * offset instance <p/> is earlier than the Initial_Instance of dataset. 1019 * @throws Exception 1020 */ 1021 private static String coord_offset_sync(int n, String timeUnit) throws Exception { 1022 Calendar rawCal = resolveOffsetRawTime(n, TimeUnit.valueOf(timeUnit), null); 1023 if (rawCal == null) { 1024 // warning already logged by resolveOffsetRawTime() 1025 return ""; 1026 } 1027 1028 int freq = getDSFrequency(); 1029 TimeUnit freqUnit = getDSTimeUnit(); 1030 int freqCount = 0; 1031 // We're going to manually turn back/forward cal by decrements/increments of freq and then check that it gives the same 1032 // time as rawCal; this is to check that the offset time resolves to a frequency offset of the effective nominal time 1033 // In other words, that there exists an integer x, such that coord:offset(n, timeUnit) == coord:current(x) is true 1034 // If not, then we'll "rewind" rawCal to the latest instance earlier than rawCal and use that. 1035 Calendar cal = getInitialInstanceCal(); 1036 if (rawCal.before(cal)) { 1037 while (cal.after(rawCal)) { 1038 cal.add(freqUnit.getCalendarUnit(), -freq); 1039 freqCount--; 1040 } 1041 } 1042 else if (rawCal.after(cal)) { 1043 while (cal.before(rawCal)) { 1044 cal.add(freqUnit.getCalendarUnit(), freq); 1045 freqCount++; 1046 } 1047 } 1048 if (cal.before(rawCal)) { 1049 rawCal = cal; 1050 } 1051 else if (cal.after(rawCal)) { 1052 cal.add(freqUnit.getCalendarUnit(), -freq); 1053 rawCal = cal; 1054 freqCount--; 1055 } 1056 String rawCalStr = DateUtils.formatDateOozieTZ(rawCal); 1057 1058 Calendar nominalInstanceCal = getInitialInstanceCal(); 1059 nominalInstanceCal.add(freqUnit.getCalendarUnit(), freq * freqCount); 1060 if (nominalInstanceCal.getTime().compareTo(getInitialInstance()) < 0) { 1061 XLog.getLog(CoordELFunctions.class).warn("If the initial instance of the dataset is later than the offset instance" 1062 + " specified, such as coord:offset({0}, {1}) in this case, an empty string is returned. This means that no" 1063 + " data is available at the offset instance specified by the user and the user could try modifying his" 1064 + " initial-instance to an earlier time.", n, timeUnit); 1065 return ""; 1066 } 1067 String nominalCalStr = DateUtils.formatDateOozieTZ(nominalInstanceCal); 1068 1069 if (!rawCalStr.equals(nominalCalStr)) { 1070 throw new RuntimeException("Shouldn't happen"); 1071 } 1072 return rawCalStr; 1073 } 1074 1075 /** 1076 * @param offset 1077 * @return n-th available latest instance Date-Time for SYNC data-set 1078 * @throws Exception 1079 */ 1080 private static String coord_latest_sync(int offset) throws Exception { 1081 return coord_latestRange_sync(offset, offset); 1082 } 1083 1084 private static String coord_latestRange_sync(int startOffset, int endOffset) throws Exception { 1085 final XLog LOG = XLog.getLog(CoordELFunctions.class); 1086 final Thread currentThread = Thread.currentThread(); 1087 ELEvaluator eval = ELEvaluator.getCurrent(); 1088 String retVal = ""; 1089 int datasetFrequency = (int) getDSFrequency();// in minutes 1090 TimeUnit dsTimeUnit = getDSTimeUnit(); 1091 int[] instCount = new int[1]; 1092 boolean useCurrentTime = Services.get().getConf().getBoolean(LATEST_EL_USE_CURRENT_TIME, false); 1093 Calendar nominalInstanceCal; 1094 if (useCurrentTime) { 1095 nominalInstanceCal = getCurrentInstance(new Date(), instCount); 1096 } 1097 else { 1098 nominalInstanceCal = getCurrentInstance(getActualTime(), instCount); 1099 } 1100 StringBuilder resolvedInstances = new StringBuilder(); 1101 StringBuilder resolvedURIPaths = new StringBuilder(); 1102 if (nominalInstanceCal != null) { 1103 Calendar initInstance = getInitialInstanceCal(); 1104 SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET); 1105 if (ds == null) { 1106 throw new RuntimeException("Associated Dataset should be defined with key " + DATASET); 1107 } 1108 String uriTemplate = ds.getUriTemplate(); 1109 Configuration conf = (Configuration) eval.getVariable(CONFIGURATION); 1110 if (conf == null) { 1111 throw new RuntimeException("Associated Configuration should be defined with key " + CONFIGURATION); 1112 } 1113 int available = 0; 1114 boolean resolved = false; 1115 String user = ParamChecker 1116 .notEmpty((String) eval.getVariable(OozieClient.USER_NAME), OozieClient.USER_NAME); 1117 String doneFlag = ds.getDoneFlag(); 1118 URIHandlerService uriService = Services.get().get(URIHandlerService.class); 1119 URIHandler uriHandler = null; 1120 Context uriContext = null; 1121 try { 1122 while (nominalInstanceCal.compareTo(initInstance) >= 0 && !currentThread.isInterrupted()) { 1123 ELEvaluator uriEval = getUriEvaluator(nominalInstanceCal); 1124 String uriPath = uriEval.evaluate(uriTemplate, String.class); 1125 if (uriHandler == null) { 1126 URI uri = new URI(uriPath); 1127 uriHandler = uriService.getURIHandler(uri); 1128 uriContext = uriHandler.getContext(uri, conf, user); 1129 } 1130 String uriWithDoneFlag = uriHandler.getURIWithDoneFlag(uriPath, doneFlag); 1131 if (uriHandler.exists(new URI(uriWithDoneFlag), uriContext)) { 1132 XLog.getLog(CoordELFunctions.class) 1133 .debug("Found latest(" + available + "): " + uriWithDoneFlag); 1134 if (available == startOffset) { 1135 LOG.debug("Matched latest(" + available + "): " + uriWithDoneFlag); 1136 resolved = true; 1137 resolvedInstances.append(DateUtils.formatDateOozieTZ(nominalInstanceCal)); 1138 resolvedURIPaths.append(uriPath); 1139 retVal = resolvedInstances.toString(); 1140 eval.setVariable("resolved_path", resolvedURIPaths.toString()); 1141 break; 1142 } 1143 else if (available <= endOffset) { 1144 LOG.debug("Matched latest(" + available + "): " + uriWithDoneFlag); 1145 resolvedInstances.append(DateUtils.formatDateOozieTZ(nominalInstanceCal)).append( 1146 INSTANCE_SEPARATOR); 1147 resolvedURIPaths.append(uriPath).append(INSTANCE_SEPARATOR); 1148 } 1149 1150 available--; 1151 } 1152 // nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), -datasetFrequency); 1153 nominalInstanceCal = (Calendar) initInstance.clone(); 1154 instCount[0]--; 1155 nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), instCount[0] * datasetFrequency); 1156 // DateUtils.moveToEnd(nominalInstanceCal, getDSEndOfFlag()); 1157 } 1158 } 1159 finally { 1160 if (uriContext != null) { 1161 uriContext.destroy(); 1162 } 1163 } 1164 if (!resolved) { 1165 // return unchanged latest function with variable 'is_resolved' 1166 // to 'false' 1167 eval.setVariable("is_resolved", Boolean.FALSE); 1168 if (startOffset == endOffset) { 1169 retVal = "${coord:latest(" + startOffset + ")}"; 1170 } 1171 else { 1172 retVal = "${coord:latestRange(" + startOffset + "," + endOffset + ")}"; 1173 } 1174 } 1175 else { 1176 eval.setVariable("is_resolved", Boolean.TRUE); 1177 } 1178 } 1179 else {// No feasible nominal time 1180 eval.setVariable("is_resolved", Boolean.FALSE); 1181 } 1182 return retVal; 1183 } 1184 1185 /** 1186 * @param tm 1187 * @return a new Evaluator to be used for URI-template evaluation 1188 */ 1189 private static ELEvaluator getUriEvaluator(Calendar tm) { 1190 tm.setTimeZone(DateUtils.getOozieProcessingTimeZone()); 1191 ELEvaluator retEval = new ELEvaluator(); 1192 retEval.setVariable("YEAR", tm.get(Calendar.YEAR)); 1193 retEval.setVariable("MONTH", (tm.get(Calendar.MONTH) + 1) < 10 ? "0" + (tm.get(Calendar.MONTH) + 1) : (tm 1194 .get(Calendar.MONTH) + 1)); 1195 retEval.setVariable("DAY", tm.get(Calendar.DAY_OF_MONTH) < 10 ? "0" + tm.get(Calendar.DAY_OF_MONTH) : tm 1196 .get(Calendar.DAY_OF_MONTH)); 1197 retEval.setVariable("HOUR", tm.get(Calendar.HOUR_OF_DAY) < 10 ? "0" + tm.get(Calendar.HOUR_OF_DAY) : tm 1198 .get(Calendar.HOUR_OF_DAY)); 1199 retEval.setVariable("MINUTE", tm.get(Calendar.MINUTE) < 10 ? "0" + tm.get(Calendar.MINUTE) : tm 1200 .get(Calendar.MINUTE)); 1201 return retEval; 1202 } 1203 1204 /** 1205 * @return whether a data set is SYNCH or ASYNC 1206 */ 1207 private static boolean isSyncDataSet() { 1208 ELEvaluator eval = ELEvaluator.getCurrent(); 1209 SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET); 1210 if (ds == null) { 1211 throw new RuntimeException("Associated Dataset should be defined with key " + DATASET); 1212 } 1213 return ds.getType().equalsIgnoreCase("SYNC"); 1214 } 1215 1216 /** 1217 * Check whether a function should be resolved. 1218 * 1219 * @param functionName 1220 * @param n 1221 * @return null if the functionName needs to be resolved otherwise return the calling function unresolved. 1222 */ 1223 private static String checkIfResolved(String functionName, String n) { 1224 ELEvaluator eval = ELEvaluator.getCurrent(); 1225 String replace = (String) eval.getVariable("resolve_" + functionName); 1226 if (replace == null || (replace != null && replace.equalsIgnoreCase("false"))) { // Don't 1227 // resolve 1228 // return "${coord:" + functionName + "(" + n +")}"; //Unresolved 1229 eval.setVariable(".wrap", "true"); 1230 return "coord:" + functionName + "(" + n + ")"; // Unresolved 1231 } 1232 return null; // Resolved it 1233 } 1234 1235 private static String echoUnResolved(String functionName, String n) { 1236 return echoUnResolvedPre(functionName, n, "coord:"); 1237 } 1238 1239 private static String echoUnResolvedPre(String functionName, String n, String prefix) { 1240 ELEvaluator eval = ELEvaluator.getCurrent(); 1241 eval.setVariable(".wrap", "true"); 1242 return prefix + functionName + "(" + n + ")"; // Unresolved 1243 } 1244 1245 /** 1246 * @return the initial instance of a DataSet in DATE 1247 */ 1248 private static Date getInitialInstance() { 1249 ELEvaluator eval = ELEvaluator.getCurrent(); 1250 return getInitialInstance(eval); 1251 } 1252 1253 /** 1254 * @return the initial instance of a DataSet in DATE 1255 */ 1256 private static Date getInitialInstance(ELEvaluator eval) { 1257 return getInitialInstanceCal(eval).getTime(); 1258 // return ds.getInitInstance(); 1259 } 1260 1261 /** 1262 * @return the initial instance of a DataSet in Calendar 1263 */ 1264 private static Calendar getInitialInstanceCal() { 1265 ELEvaluator eval = ELEvaluator.getCurrent(); 1266 return getInitialInstanceCal(eval); 1267 } 1268 1269 /** 1270 * @return the initial instance of a DataSet in Calendar 1271 */ 1272 private static Calendar getInitialInstanceCal(ELEvaluator eval) { 1273 SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET); 1274 if (ds == null) { 1275 throw new RuntimeException("Associated Dataset should be defined with key " + DATASET); 1276 } 1277 Calendar effInitTS = new GregorianCalendar(ds.getTimeZone()); 1278 effInitTS.setTime(ds.getInitInstance()); 1279 // To adjust EOD/EOM 1280 DateUtils.moveToEnd(effInitTS, getDSEndOfFlag(eval)); 1281 return effInitTS; 1282 // return ds.getInitInstance(); 1283 } 1284 1285 /** 1286 * @return Nominal or action creation Time when all the dependencies of an application instance are met. 1287 */ 1288 private static Date getActionCreationtime() { 1289 ELEvaluator eval = ELEvaluator.getCurrent(); 1290 return getActionCreationtime(eval); 1291 } 1292 1293 /** 1294 * @return Nominal or action creation Time when all the dependencies of an application instance are met. 1295 */ 1296 private static Date getActionCreationtime(ELEvaluator eval) { 1297 SyncCoordAction coordAction = (SyncCoordAction) eval.getVariable(COORD_ACTION); 1298 if (coordAction == null) { 1299 throw new RuntimeException("Associated Application instance should be defined with key " + COORD_ACTION); 1300 } 1301 return coordAction.getNominalTime(); 1302 } 1303 1304 /** 1305 * @return Actual Time when all the dependencies of an application instance are met. 1306 */ 1307 private static Date getActualTime() { 1308 ELEvaluator eval = ELEvaluator.getCurrent(); 1309 SyncCoordAction coordAction = (SyncCoordAction) eval.getVariable(COORD_ACTION); 1310 if (coordAction == null) { 1311 throw new RuntimeException("Associated Application instance should be defined with key " + COORD_ACTION); 1312 } 1313 return coordAction.getActualTime(); 1314 } 1315 1316 /** 1317 * @return TimeZone for the application or job. 1318 */ 1319 private static TimeZone getJobTZ() { 1320 ELEvaluator eval = ELEvaluator.getCurrent(); 1321 SyncCoordAction coordAction = (SyncCoordAction) eval.getVariable(COORD_ACTION); 1322 if (coordAction == null) { 1323 throw new RuntimeException("Associated Application instance should be defined with key " + COORD_ACTION); 1324 } 1325 return coordAction.getTimeZone(); 1326 } 1327 1328 /** 1329 * Find the current instance based on effectiveTime (i.e Action_Creation_Time or Action_Start_Time) 1330 * 1331 * @return current instance i.e. current(0) returns null if effectiveTime is earlier than Initial Instance time of 1332 * the dataset. 1333 */ 1334 public static Calendar getCurrentInstance(Date effectiveTime, int instanceCount[]) { 1335 ELEvaluator eval = ELEvaluator.getCurrent(); 1336 return getCurrentInstance(effectiveTime, instanceCount, eval); 1337 } 1338 1339 /** 1340 * Find the current instance based on effectiveTime (i.e Action_Creation_Time or Action_Start_Time) 1341 * 1342 * @return current instance i.e. current(0) returns null if effectiveTime is earlier than Initial Instance time of 1343 * the dataset. 1344 */ 1345 private static Calendar getCurrentInstance(Date effectiveTime, int instanceCount[], ELEvaluator eval) { 1346 Date datasetInitialInstance = getInitialInstance(eval); 1347 TimeUnit dsTimeUnit = getDSTimeUnit(eval); 1348 TimeZone dsTZ = getDatasetTZ(eval); 1349 int dsFreq = getDSFrequency(eval); 1350 // Convert Date to Calendar for corresponding TZ 1351 Calendar current = Calendar.getInstance(dsTZ); 1352 current.setTime(datasetInitialInstance); 1353 1354 Calendar calEffectiveTime = new GregorianCalendar(dsTZ); 1355 calEffectiveTime.setTime(effectiveTime); 1356 if (instanceCount == null) { // caller doesn't care about this value 1357 instanceCount = new int[1]; 1358 } 1359 instanceCount[0] = 0; 1360 if (current.compareTo(calEffectiveTime) > 0) { 1361 return null; 1362 } 1363 1364 switch(dsTimeUnit) { 1365 case MINUTE: 1366 instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / MINUTE_MSEC); 1367 break; 1368 case HOUR: 1369 instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / HOUR_MSEC); 1370 break; 1371 case DAY: 1372 case END_OF_DAY: 1373 instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / DAY_MSEC); 1374 break; 1375 case MONTH: 1376 case END_OF_MONTH: 1377 instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / MONTH_MSEC); 1378 break; 1379 case YEAR: 1380 instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / YEAR_MSEC); 1381 break; 1382 default: 1383 throw new IllegalArgumentException("Unhandled dataset time unit " + dsTimeUnit); 1384 } 1385 1386 if (instanceCount[0] > 2) { 1387 instanceCount[0] = (instanceCount[0] / dsFreq); 1388 current.add(dsTimeUnit.getCalendarUnit(), instanceCount[0] * dsFreq); 1389 } else { 1390 instanceCount[0] = 0; 1391 } 1392 while (!current.getTime().after(effectiveTime)) { 1393 current.add(dsTimeUnit.getCalendarUnit(), dsFreq); 1394 instanceCount[0]++; 1395 } 1396 current.add(dsTimeUnit.getCalendarUnit(), -dsFreq); 1397 instanceCount[0]--; 1398 return current; 1399 } 1400 1401 /** 1402 * Find the current instance based on effectiveTime (i.e Action_Creation_Time or Action_Start_Time) 1403 * 1404 * @return current instance i.e. current(0) returns null if effectiveTime is earlier than Initial Instance time of 1405 * the dataset. 1406 */ 1407 private static Calendar getCurrentInstance_old(Date effectiveTime, int instanceCount[], ELEvaluator eval) { 1408 Date datasetInitialInstance = getInitialInstance(eval); 1409 TimeUnit dsTimeUnit = getDSTimeUnit(eval); 1410 TimeZone dsTZ = getDatasetTZ(eval); 1411 int dsFreq = getDSFrequency(eval); 1412 // Convert Date to Calendar for corresponding TZ 1413 Calendar current = Calendar.getInstance(); 1414 current.setTime(datasetInitialInstance); 1415 current.setTimeZone(dsTZ); 1416 1417 Calendar calEffectiveTime = Calendar.getInstance(); 1418 calEffectiveTime.setTime(effectiveTime); 1419 calEffectiveTime.setTimeZone(dsTZ); 1420 if (instanceCount == null) { // caller doesn't care about this value 1421 instanceCount = new int[1]; 1422 } 1423 instanceCount[0] = 0; 1424 if (current.compareTo(calEffectiveTime) > 0) { 1425 return null; 1426 } 1427 Calendar origCurrent = (Calendar) current.clone(); 1428 while (current.compareTo(calEffectiveTime) <= 0) { 1429 current = (Calendar) origCurrent.clone(); 1430 instanceCount[0]++; 1431 current.add(dsTimeUnit.getCalendarUnit(), instanceCount[0] * dsFreq); 1432 } 1433 instanceCount[0]--; 1434 1435 current = (Calendar) origCurrent.clone(); 1436 current.add(dsTimeUnit.getCalendarUnit(), instanceCount[0] * dsFreq); 1437 return current; 1438 } 1439 1440 public static Calendar getEffectiveNominalTime() { 1441 Date datasetInitialInstance = getInitialInstance(); 1442 TimeZone dsTZ = getDatasetTZ(); 1443 // Convert Date to Calendar for corresponding TZ 1444 Calendar current = Calendar.getInstance(); 1445 current.setTime(datasetInitialInstance); 1446 current.setTimeZone(dsTZ); 1447 1448 Calendar calEffectiveTime = Calendar.getInstance(); 1449 calEffectiveTime.setTime(getActionCreationtime()); 1450 calEffectiveTime.setTimeZone(dsTZ); 1451 if (current.compareTo(calEffectiveTime) > 0) { 1452 // Nominal Time < initial Instance 1453 // TODO: getClass() call doesn't work from static method. 1454 // XLog.getLog("CoordELFunction.class").warn("ACTION CREATED BEFORE INITIAL INSTACE "+ 1455 // current.getTime()); 1456 return null; 1457 } 1458 return calEffectiveTime; 1459 } 1460 1461 /** 1462 * @return dataset frequency in minutes 1463 */ 1464 private static int getDSFrequency() { 1465 ELEvaluator eval = ELEvaluator.getCurrent(); 1466 return getDSFrequency(eval); 1467 } 1468 1469 /** 1470 * @return dataset frequency in minutes 1471 */ 1472 private static int getDSFrequency(ELEvaluator eval) { 1473 SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET); 1474 if (ds == null) { 1475 throw new RuntimeException("Associated Dataset should be defined with key " + DATASET); 1476 } 1477 return ds.getFrequency(); 1478 } 1479 1480 /** 1481 * @return dataset TimeUnit 1482 */ 1483 private static TimeUnit getDSTimeUnit() { 1484 ELEvaluator eval = ELEvaluator.getCurrent(); 1485 return getDSTimeUnit(eval); 1486 } 1487 1488 /** 1489 * @return dataset TimeUnit 1490 */ 1491 public static TimeUnit getDSTimeUnit(ELEvaluator eval) { 1492 SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET); 1493 if (ds == null) { 1494 throw new RuntimeException("Associated Dataset should be defined with key " + DATASET); 1495 } 1496 return ds.getTimeUnit(); 1497 } 1498 1499 /** 1500 * @return dataset TimeZone 1501 */ 1502 public static TimeZone getDatasetTZ() { 1503 ELEvaluator eval = ELEvaluator.getCurrent(); 1504 return getDatasetTZ(eval); 1505 } 1506 1507 /** 1508 * @return dataset TimeZone 1509 */ 1510 private static TimeZone getDatasetTZ(ELEvaluator eval) { 1511 SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET); 1512 if (ds == null) { 1513 throw new RuntimeException("Associated Dataset should be defined with key " + DATASET); 1514 } 1515 return ds.getTimeZone(); 1516 } 1517 1518 /** 1519 * @return dataset TimeUnit 1520 */ 1521 private static TimeUnit getDSEndOfFlag() { 1522 ELEvaluator eval = ELEvaluator.getCurrent(); 1523 return getDSEndOfFlag(eval); 1524 } 1525 1526 /** 1527 * @return dataset TimeUnit 1528 */ 1529 private static TimeUnit getDSEndOfFlag(ELEvaluator eval) { 1530 SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET); 1531 if (ds == null) { 1532 throw new RuntimeException("Associated Dataset should be defined with key " + DATASET); 1533 } 1534 return ds.getEndOfDuration();// == null ? "": ds.getEndOfDuration(); 1535 } 1536 1537 /** 1538 * Return a job configuration property for the coordinator. 1539 * 1540 * @param property property name. 1541 * @return the value of the property, <code>null</code> if the property is undefined. 1542 */ 1543 public static String coord_conf(String property) { 1544 ELEvaluator eval = ELEvaluator.getCurrent(); 1545 return (String) eval.getVariable(property); 1546 } 1547 1548 /** 1549 * Return the user that submitted the coordinator job. 1550 * 1551 * @return the user that submitted the coordinator job. 1552 */ 1553 public static String coord_user() { 1554 ELEvaluator eval = ELEvaluator.getCurrent(); 1555 return (String) eval.getVariable(OozieClient.USER_NAME); 1556 } 1557 1558 /** 1559 * Takes two offset times and returns a list of multiples of the frequency offset from the effective nominal time that occur 1560 * between them. The caller should make sure that startCal is earlier than endCal. 1561 * <p> 1562 * As a simple example, assume its the same day: startCal is 1:00, endCal is 2:00, frequency is 20min, and effective nominal 1563 * time is 1:20 -- then this method would return a list containing: -20, 0, 20, 40, 60 1564 * 1565 * @param startCal The earlier offset time 1566 * @param endCal The later offset time 1567 * @param eval The ELEvaluator to use; cannot be null 1568 * @return A list of multiple of the frequency offset from the effective nominal time that occur between the startCal and endCal 1569 */ 1570 public static List<Integer> expandOffsetTimes(Calendar startCal, Calendar endCal, ELEvaluator eval) { 1571 List<Integer> expandedFreqs = new ArrayList<Integer>(); 1572 // Use eval because the "current" eval isn't set 1573 int freq = getDSFrequency(eval); 1574 TimeUnit freqUnit = getDSTimeUnit(eval); 1575 Calendar cal = getCurrentInstance(getActionCreationtime(eval), null, eval); 1576 int totalFreq = 0; 1577 if (startCal.before(cal)) { 1578 while (cal.after(startCal)) { 1579 cal.add(freqUnit.getCalendarUnit(), -freq); 1580 totalFreq += -freq; 1581 } 1582 if (cal.before(startCal)) { 1583 cal.add(freqUnit.getCalendarUnit(), freq); 1584 totalFreq += freq; 1585 } 1586 } 1587 else if (startCal.after(cal)) { 1588 while (cal.before(startCal)) { 1589 cal.add(freqUnit.getCalendarUnit(), freq); 1590 totalFreq += freq; 1591 } 1592 } 1593 // At this point, cal is the smallest multiple of the dataset frequency that is >= to the startCal and offset from the 1594 // effective nominal time. Now we can find all of the instances that occur between startCal and endCal, inclusive. 1595 while (cal.before(endCal) || cal.equals(endCal)) { 1596 expandedFreqs.add(totalFreq); 1597 cal.add(freqUnit.getCalendarUnit(), freq); 1598 totalFreq += freq; 1599 } 1600 return expandedFreqs; 1601 } 1602 1603 /** 1604 * Resolve the offset time from the effective nominal time 1605 * 1606 * @param n offset amount (integer) 1607 * @param timeUnit TimeUnit for offset n ("MINUTE", "HOUR", "DAY", "MONTH", "YEAR") 1608 * @param eval The ELEvaluator to use; or null to use the "current" eval 1609 * @return A Calendar of the offset time 1610 */ 1611 public static Calendar resolveOffsetRawTime(int n, TimeUnit timeUnit, ELEvaluator eval) { 1612 // Use eval if given (for when the "current" eval isn't set) 1613 Calendar cal; 1614 if (eval == null) { 1615 cal = getCurrentInstance(getActionCreationtime(), null); 1616 } 1617 else { 1618 cal = getCurrentInstance(getActionCreationtime(eval), null, eval); 1619 } 1620 if (cal == null) { 1621 XLog.getLog(CoordELFunctions.class).warn("If the initial instance of the dataset is later than the nominal time, an" 1622 + " empty string is returned. This means that no data is available at the offset instance specified by the user" 1623 + " and the user could try modifying his or her initial-instance to an earlier time."); 1624 return null; 1625 } 1626 cal.add(timeUnit.getCalendarUnit(), n); 1627 return cal; 1628 } 1629}