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.executor.jpa; 020 021import java.sql.Timestamp; 022import java.text.ParseException; 023import java.util.ArrayList; 024import java.util.Collections; 025import java.util.HashMap; 026import java.util.Iterator; 027import java.util.List; 028import java.util.Map; 029import java.util.Map.Entry; 030import java.util.regex.Matcher; 031import java.util.regex.Pattern; 032 033import javax.persistence.EntityManager; 034import javax.persistence.Query; 035 036import org.apache.oozie.BulkResponseInfo; 037import org.apache.oozie.BundleJobBean; 038import org.apache.oozie.CoordinatorActionBean; 039import org.apache.oozie.CoordinatorJobBean; 040import org.apache.oozie.ErrorCode; 041import org.apache.oozie.StringBlob; 042import org.apache.oozie.client.BundleJob; 043import org.apache.oozie.client.CoordinatorAction; 044import org.apache.oozie.client.CoordinatorJob; 045import org.apache.oozie.client.rest.BulkResponseImpl; 046import org.apache.oozie.service.Services; 047import org.apache.oozie.util.DateUtils; 048import org.apache.oozie.util.ParamChecker; 049 050/** 051 * The query executor class for bulk monitoring queries i.e. debugging bundle -> 052 * coord actions directly 053 */ 054public class BulkJPAExecutor implements JPAExecutor<BulkResponseInfo> { 055 private Map<String, List<String>> bulkFilter; 056 // defaults 057 private int start = 1; 058 private int len = 50; 059 private enum PARAM_TYPE { 060 ID, NAME 061 } 062 063 public BulkJPAExecutor(Map<String, List<String>> bulkFilter, int start, int len) { 064 ParamChecker.notNull(bulkFilter, "bulkFilter"); 065 this.bulkFilter = bulkFilter; 066 this.start = start; 067 this.len = len; 068 } 069 070 @Override 071 public String getName() { 072 return "BulkJPAExecutor"; 073 } 074 075 @Override 076 public BulkResponseInfo execute(EntityManager em) throws JPAExecutorException { 077 List<BulkResponseImpl> responseList = new ArrayList<BulkResponseImpl>(); 078 Map<String, Timestamp> actionTimes = new HashMap<String, Timestamp>(); 079 080 try { 081 List<String> coords = bulkFilter.get(BulkResponseImpl.BULK_FILTER_COORD); 082 List<String> statuses = bulkFilter.get(BulkResponseImpl.BULK_FILTER_STATUS); 083 List<String> params = new ArrayList<String>(); 084 085 // Lightweight Query 1 on Bundle level to fetch the bundle job(s) 086 // corresponding to names or ids 087 List<BundleJobBean> bundleBeans = bundleQuery(em); 088 089 // Join query between coordinator job and coordinator action tables 090 // to get entries for specific bundleId only 091 String conditions = actionQuery(coords, params, statuses, em, bundleBeans, actionTimes, responseList); 092 093 // Query to get the count of records 094 long total = countQuery(statuses, params, conditions, em, bundleBeans, actionTimes); 095 096 BulkResponseInfo bulk = new BulkResponseInfo(responseList, start, len, total); 097 return bulk; 098 } 099 catch (Exception e) { 100 throw new JPAExecutorException(ErrorCode.E0603, e.getMessage(), e); 101 } 102 } 103 104 /** 105 * build the bundle level query to get bundle beans for the specified ids or appnames 106 * @param em 107 * @return List BundleJobBeans 108 * @throws JPAExecutorException 109 */ 110 @SuppressWarnings("unchecked") 111 private List<BundleJobBean> bundleQuery(EntityManager em) throws JPAExecutorException { 112 Query q = em.createNamedQuery("BULK_MONITOR_BUNDLE_QUERY"); 113 StringBuilder bundleQuery = new StringBuilder(q.toString()); 114 115 StringBuilder whereClause = null; 116 List<String> bundles = bulkFilter.get(BulkResponseImpl.BULK_FILTER_BUNDLE); 117 if (bundles != null) { 118 PARAM_TYPE type = getParamType(bundles.get(0), 'B'); 119 if (type == PARAM_TYPE.NAME) { 120 whereClause = inClause(bundles.size(), "appName", 'b', "bundles"); 121 } 122 else if (type == PARAM_TYPE.ID) { 123 whereClause = inClause(bundles.size(), "id", 'b', "bundles"); 124 } 125 126 // Query: select <columns> from BundleJobBean b where b.id IN (...) _or_ b.appName IN (...) 127 bundleQuery.append(whereClause.replace(whereClause.indexOf("AND"), whereClause.indexOf("AND") + 3, "WHERE")); 128 Query tmp = em.createQuery(bundleQuery.toString()); 129 130 fillParameters(tmp, "bundles", bundles); 131 132 List<Object[]> bundleObjs = (List<Object[]>) tmp.getResultList(); 133 if (bundleObjs.isEmpty()) { 134 throw new JPAExecutorException(ErrorCode.E0603, "No entries found for given bundle(s)"); 135 } 136 137 List<BundleJobBean> bundleBeans = new ArrayList<BundleJobBean>(); 138 for (Object[] bundleElem : bundleObjs) { 139 bundleBeans.add(constructBundleBean(bundleElem)); 140 } 141 return bundleBeans; 142 } 143 return null; 144 } 145 146 /** 147 * Validate and determine whether passed param is job-id or appname 148 * @param id 149 * @param job 150 * @return PARAM_TYPE 151 */ 152 private PARAM_TYPE getParamType(String id, char job) { 153 Pattern p = Pattern.compile("\\d{7}-\\d{15}-" + Services.get().getSystemId() + "-" + job); 154 Matcher m = p.matcher(id); 155 if (m.matches()) { 156 return PARAM_TYPE.ID; 157 } 158 return PARAM_TYPE.NAME; 159 } 160 161 /** 162 * Compose the coord action level query comprising bundle id/appname filter and coord action 163 * status filter (if specified) and start-time or nominal-time filter (if specified) 164 * @param em 165 * @param bundles 166 * @param times 167 * @param responseList 168 * @return Query string 169 * @throws ParseException 170 */ 171 @SuppressWarnings("unchecked") 172 private String actionQuery(final List<String> coords, final List<String> params, List<String> statuses, EntityManager em, 173 List<BundleJobBean> bundles, Map<String, Timestamp> times, List<BulkResponseImpl> responseList) 174 throws ParseException { 175 Query q = em.createNamedQuery("BULK_MONITOR_ACTIONS_QUERY"); 176 StringBuilder getActions = new StringBuilder(q.toString()); 177 int offset = getActions.indexOf("ORDER"); 178 StringBuilder conditionClause = new StringBuilder(); 179 180 // Query: Select <columns> from CoordinatorActionBean a, CoordinatorJobBean c WHERE a.jobId = c.id 181 // AND c.bundleId = :bundleId AND c.appName/id IN (...) 182 183 if (coords != null) { 184 PARAM_TYPE type = getParamType(coords.get(0), 'C'); 185 if (type == PARAM_TYPE.NAME) { 186 conditionClause.append(inClause(coords.size(), "appName", 'c', "param")); 187 params.addAll(coords); 188 } 189 else if (type == PARAM_TYPE.ID) { 190 conditionClause.append(inClause(coords.size(), "id", 'c', "param")); 191 params.addAll(coords); 192 } 193 } 194 // Query: Select <columns> from CoordinatorActionBean a, CoordinatorJobBean c WHERE a.jobId = c.id 195 // AND c.bundleId = :bundleId AND c.appName/id IN (...) AND a.statusStr IN (...) 196 conditionClause.append(statusClause(statuses)); 197 198 offset = getActions.indexOf("ORDER"); 199 getActions.insert(offset - 1, conditionClause); 200 201 // Query: Select <columns> from CoordinatorActionBean a, CoordinatorJobBean c WHERE a.jobId = c.id 202 // AND c.bundleId = :bundleId AND c.appName/id IN (...) AND a.statusStr IN (...) 203 // AND a.createdTimestamp >= startCreated _or_ a.createdTimestamp <= endCreated 204 // AND a.nominalTimestamp >= startNominal _or_ a.nominalTimestamp <= endNominal 205 timesClause(getActions, offset, times); 206 q = em.createQuery(getActions.toString()); 207 Iterator<Entry<String, Timestamp>> iter = times.entrySet().iterator(); 208 while (iter.hasNext()) { 209 Entry<String, Timestamp> time = iter.next(); 210 q.setParameter(time.getKey(), time.getValue()); 211 } 212 // pagination 213 q.setFirstResult(start - 1); 214 q.setMaxResults(len); 215 216 if (coords != null) { 217 fillParameters(q, "param", coords); 218 } 219 220 if (statuses != null) { 221 fillParameters(q, "status", statuses); 222 } 223 224 // repeatedly execute above query for each bundle 225 for (BundleJobBean bundle : bundles) { 226 q.setParameter("bundleId", bundle.getId()); 227 List<Object[]> response = q.getResultList(); 228 for (Object[] r : response) { 229 BulkResponseImpl br = getResponseFromObject(bundle, r); 230 responseList.add(br); 231 } 232 } 233 return q.toString(); 234 } 235 236 /** 237 * Get total number of records for use with offset and len in API 238 * @param clause 239 * @param em 240 * @param bundles 241 * @return total count of coord actions 242 */ 243 private long countQuery(List<String> statuses, List<String> params, String clause, EntityManager em, 244 List<BundleJobBean> bundles, Map<String, Timestamp> times) { 245 Query q = em.createNamedQuery("BULK_MONITOR_COUNT_QUERY"); 246 StringBuilder getTotal = new StringBuilder(q.toString() + " "); 247 // Query: select COUNT(a) from CoordinatorActionBean a, CoordinatorJobBean c 248 // get entire WHERE clause from above i.e. actionQuery() for all conditions on coordinator job 249 // and action status and times 250 getTotal.append(clause.substring(clause.indexOf("WHERE"), clause.indexOf("ORDER"))); 251 int offset = getTotal.indexOf("bundleId"); 252 List<String> bundleIds = new ArrayList<String>(); 253 for (BundleJobBean bundle : bundles) { 254 bundleIds.add(bundle.getId()); 255 } 256 // Query: select COUNT(a) from CoordinatorActionBean a, CoordinatorJobBean c WHERE ... 257 // AND c.bundleId IN (... list of bundle ids) i.e. replace single :bundleId with list 258 getTotal = getTotal.replace(offset - 6, offset + 20, inClause(bundleIds.size(), "bundleId", 'c', "count").toString()); 259 q = em.createQuery(getTotal.toString()); 260 261 fillParameters(q, "count", bundleIds); 262 263 if (statuses != null) { 264 fillParameters(q, "status", statuses); 265 } 266 267 if (params != null) { 268 fillParameters(q, "param", params); 269 } 270 271 Iterator<Entry<String, Timestamp>> iter = times.entrySet().iterator(); 272 while (iter.hasNext()) { 273 Entry<String, Timestamp> time = iter.next(); 274 q.setParameter(time.getKey(), time.getValue()); 275 } 276 long total = ((Long) q.getSingleResult()).longValue(); 277 return total; 278 } 279 280 // Form the where clause to filter by coordinator appname/id 281 private StringBuilder inClause(int noOfValues, String col, char type, String paramPrefix) { 282 StringBuilder sb = new StringBuilder(); 283 boolean firstVal = true; 284 285 for (int i = 0; i < noOfValues; i++) { 286 if (firstVal) { 287 sb.append(" AND " + type + "." + col + " IN (:" + paramPrefix + i); 288 firstVal = false; 289 } 290 else { 291 sb.append(", :" + paramPrefix + i); 292 } 293 } 294 if (!firstVal) { 295 sb.append(") "); 296 } 297 298 return sb; 299 } 300 301 // Form the where clause to filter by coord action status 302 private StringBuilder statusClause(List<String> statuses) { 303 StringBuilder sb = new StringBuilder(); 304 305 if (statuses != null) { 306 sb = inClause(statuses.size(), "statusStr", 'a', "status"); 307 } 308 309 if (sb.length() == 0) { // statuses was null. adding default 310 sb.append(" AND a.statusStr IN ('KILLED', 'FAILED') "); 311 } 312 313 return sb; 314 } 315 316 private void timesClause(StringBuilder sb, int offset, Map<String, Timestamp> eachTime) throws ParseException { 317 Timestamp ts = null; 318 List<String> times = bulkFilter.get(BulkResponseImpl.BULK_FILTER_START_CREATED_EPOCH); 319 if (times != null) { 320 ts = new Timestamp(DateUtils.parseDateUTC(times.get(0)).getTime()); 321 sb.insert(offset - 1, " AND a.createdTimestamp >= :startCreated"); 322 eachTime.put("startCreated", ts); 323 } 324 times = bulkFilter.get(BulkResponseImpl.BULK_FILTER_END_CREATED_EPOCH); 325 if (times != null) { 326 ts = new Timestamp(DateUtils.parseDateUTC(times.get(0)).getTime()); 327 sb.insert(offset - 1, " AND a.createdTimestamp <= :endCreated"); 328 eachTime.put("endCreated", ts); 329 } 330 times = bulkFilter.get(BulkResponseImpl.BULK_FILTER_START_NOMINAL_EPOCH); 331 if (times != null) { 332 ts = new Timestamp(DateUtils.parseDateUTC(times.get(0)).getTime()); 333 sb.insert(offset - 1, " AND a.nominalTimestamp >= :startNominal"); 334 eachTime.put("startNominal", ts); 335 } 336 times = bulkFilter.get(BulkResponseImpl.BULK_FILTER_END_NOMINAL_EPOCH); 337 if (times != null) { 338 ts = new Timestamp(DateUtils.parseDateUTC(times.get(0)).getTime()); 339 sb.insert(offset - 1, " AND a.nominalTimestamp <= :endNominal"); 340 eachTime.put("endNominal", ts); 341 } 342 } 343 344 private BulkResponseImpl getResponseFromObject(BundleJobBean bundleBean, Object arr[]) { 345 BulkResponseImpl bean = new BulkResponseImpl(); 346 CoordinatorJobBean coordBean = new CoordinatorJobBean(); 347 CoordinatorActionBean actionBean = new CoordinatorActionBean(); 348 if (arr[0] != null) { 349 actionBean.setId((String) arr[0]); 350 } 351 if (arr[1] != null) { 352 actionBean.setActionNumber((Integer) arr[1]); 353 } 354 if (arr[2] != null) { 355 actionBean.setErrorCode((String) arr[2]); 356 } 357 if (arr[3] != null) { 358 actionBean.setErrorMessage((String) arr[3]); 359 } 360 if (arr[4] != null) { 361 actionBean.setExternalId((String) arr[4]); 362 } 363 if (arr[5] != null) { 364 actionBean.setExternalStatus((String) arr[5]); 365 } 366 if (arr[6] != null) { 367 actionBean.setStatus(CoordinatorAction.Status.valueOf((String) arr[6])); 368 } 369 if (arr[7] != null) { 370 actionBean.setCreatedTime(DateUtils.toDate((Timestamp) arr[7])); 371 } 372 if (arr[8] != null) { 373 actionBean.setNominalTime(DateUtils.toDate((Timestamp) arr[8])); 374 } 375 if (arr[9] != null) { 376 actionBean.setMissingDependenciesBlob((StringBlob) arr[9]); 377 } 378 if (arr[10] != null) { 379 coordBean.setId((String) arr[10]); 380 actionBean.setJobId((String) arr[10]); 381 } 382 if (arr[11] != null) { 383 coordBean.setAppName((String) arr[11]); 384 } 385 if (arr[12] != null) { 386 coordBean.setStatus(CoordinatorJob.Status.valueOf((String) arr[12])); 387 } 388 bean.setBundle(bundleBean); 389 bean.setCoordinator(coordBean); 390 bean.setAction(actionBean); 391 return bean; 392 } 393 394 private BundleJobBean constructBundleBean(Object[] barr) throws JPAExecutorException { 395 BundleJobBean bean = new BundleJobBean(); 396 if (barr[0] != null) { 397 bean.setId((String) barr[0]); 398 } 399 else { 400 throw new JPAExecutorException(ErrorCode.E0603, 401 "bundleId returned by query is null - cannot retrieve bulk results"); 402 } 403 if (barr[1] != null) { 404 bean.setAppName((String) barr[1]); 405 } 406 if (barr[2] != null) { 407 bean.setStatus(BundleJob.Status.valueOf((String) barr[2])); 408 } 409 if (barr[3] != null) { 410 bean.setUser((String) barr[3]); 411 } 412 return bean; 413 } 414 415 private void fillParameters(Query query, String prefix, List<?> values) { 416 for (int i = 0; i < values.size(); i++) { 417 query.setParameter(prefix + i, values.get(i)); 418 } 419 } 420 421 // null safeguard 422 public static List<String> nullToEmpty(List<String> input) { 423 return input == null ? Collections.<String> emptyList() : input; 424 } 425 426}