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 org.apache.oozie.BinaryBlob; 022import org.apache.oozie.ErrorCode; 023import org.apache.oozie.StringBlob; 024import org.apache.oozie.WorkflowJobBean; 025import org.apache.oozie.service.JPAService; 026import org.apache.oozie.service.Services; 027import org.apache.oozie.util.DateUtils; 028 029import javax.persistence.EntityManager; 030import javax.persistence.Query; 031import java.sql.Timestamp; 032import java.util.ArrayList; 033import java.util.List; 034 035/** 036 * Query Executor that provides API to run query for Workflow Job 037 */ 038public class WorkflowJobQueryExecutor extends QueryExecutor<WorkflowJobBean, WorkflowJobQueryExecutor.WorkflowJobQuery> { 039 040 public enum WorkflowJobQuery { 041 UPDATE_WORKFLOW, 042 UPDATE_WORKFLOW_MODTIME, 043 UPDATE_WORKFLOW_STATUS_MODTIME, 044 UPDATE_WORKFLOW_PARENT_MODIFIED, 045 UPDATE_WORKFLOW_STATUS_INSTANCE_MODIFIED, 046 UPDATE_WORKFLOW_STATUS_INSTANCE_MOD_END, 047 UPDATE_WORKFLOW_STATUS_INSTANCE_MOD_START_END, 048 UPDATE_WORKFLOW_RERUN, 049 GET_WORKFLOW, 050 GET_WORKFLOW_STARTTIME, 051 GET_WORKFLOW_START_END_TIME, 052 GET_WORKFLOW_USER_GROUP, 053 GET_WORKFLOW_SUSPEND, 054 GET_WORKFLOW_ACTION_OP, 055 GET_WORKFLOW_RERUN, 056 GET_WORKFLOW_DEFINITION, 057 GET_WORKFLOW_KILL, 058 GET_WORKFLOW_RESUME, 059 GET_WORKFLOW_STATUS, 060 GET_WORKFLOWS_PARENT_COORD_RERUN, 061 GET_COMPLETED_COORD_WORKFLOWS_OLDER_THAN 062 }; 063 064 private static WorkflowJobQueryExecutor instance = new WorkflowJobQueryExecutor(); 065 066 private WorkflowJobQueryExecutor() { 067 } 068 069 public static QueryExecutor<WorkflowJobBean, WorkflowJobQueryExecutor.WorkflowJobQuery> getInstance() { 070 return WorkflowJobQueryExecutor.instance; 071 } 072 073 @Override 074 public Query getUpdateQuery(WorkflowJobQuery namedQuery, WorkflowJobBean wfBean, EntityManager em) 075 throws JPAExecutorException { 076 077 Query query = em.createNamedQuery(namedQuery.name()); 078 switch (namedQuery) { 079 case UPDATE_WORKFLOW: 080 query.setParameter("appName", wfBean.getAppName()); 081 query.setParameter("appPath", wfBean.getAppPath()); 082 query.setParameter("conf", wfBean.getConfBlob()); 083 query.setParameter("groupName", wfBean.getGroup()); 084 query.setParameter("run", wfBean.getRun()); 085 query.setParameter("user", wfBean.getUser()); 086 query.setParameter("createdTime", wfBean.getCreatedTimestamp()); 087 query.setParameter("endTime", wfBean.getEndTimestamp()); 088 query.setParameter("externalId", wfBean.getExternalId()); 089 query.setParameter("lastModTime", wfBean.getLastModifiedTimestamp()); 090 query.setParameter("logToken", wfBean.getLogToken()); 091 query.setParameter("protoActionConf", wfBean.getProtoActionConfBlob()); 092 query.setParameter("slaXml", wfBean.getSlaXmlBlob()); 093 query.setParameter("startTime", wfBean.getStartTimestamp()); 094 query.setParameter("status", wfBean.getStatusStr()); 095 query.setParameter("wfInstance", wfBean.getWfInstanceBlob()); 096 query.setParameter("id", wfBean.getId()); 097 break; 098 case UPDATE_WORKFLOW_MODTIME: 099 query.setParameter("lastModTime", wfBean.getLastModifiedTimestamp()); 100 query.setParameter("id", wfBean.getId()); 101 break; 102 case UPDATE_WORKFLOW_STATUS_MODTIME: 103 query.setParameter("status", wfBean.getStatus().toString()); 104 query.setParameter("lastModTime", wfBean.getLastModifiedTimestamp()); 105 query.setParameter("id", wfBean.getId()); 106 break; 107 case UPDATE_WORKFLOW_PARENT_MODIFIED: 108 query.setParameter("parentId", wfBean.getParentId()); 109 query.setParameter("lastModTime", wfBean.getLastModifiedTimestamp()); 110 query.setParameter("id", wfBean.getId()); 111 break; 112 case UPDATE_WORKFLOW_STATUS_INSTANCE_MODIFIED: 113 query.setParameter("status", wfBean.getStatus().toString()); 114 query.setParameter("wfInstance", wfBean.getWfInstanceBlob()); 115 query.setParameter("lastModTime", wfBean.getLastModifiedTimestamp()); 116 query.setParameter("id", wfBean.getId()); 117 break; 118 case UPDATE_WORKFLOW_STATUS_INSTANCE_MOD_END: 119 query.setParameter("status", wfBean.getStatus().toString()); 120 query.setParameter("wfInstance", wfBean.getWfInstanceBlob()); 121 query.setParameter("lastModTime", wfBean.getLastModifiedTimestamp()); 122 query.setParameter("endTime", wfBean.getEndTimestamp()); 123 query.setParameter("id", wfBean.getId()); 124 break; 125 case UPDATE_WORKFLOW_STATUS_INSTANCE_MOD_START_END: 126 query.setParameter("status", wfBean.getStatus().toString()); 127 query.setParameter("wfInstance", wfBean.getWfInstanceBlob()); 128 query.setParameter("lastModTime", wfBean.getLastModifiedTimestamp()); 129 query.setParameter("startTime", wfBean.getStartTimestamp()); 130 query.setParameter("endTime", wfBean.getEndTimestamp()); 131 query.setParameter("id", wfBean.getId()); 132 break; 133 case UPDATE_WORKFLOW_RERUN: 134 query.setParameter("appName", wfBean.getAppName()); 135 query.setParameter("protoActionConf", wfBean.getProtoActionConfBlob()); 136 query.setParameter("appPath", wfBean.getAppPath()); 137 query.setParameter("conf", wfBean.getConfBlob()); 138 query.setParameter("logToken", wfBean.getLogToken()); 139 query.setParameter("user", wfBean.getUser()); 140 query.setParameter("group", wfBean.getGroup()); 141 query.setParameter("externalId", wfBean.getExternalId()); 142 query.setParameter("endTime", wfBean.getEndTimestamp()); 143 query.setParameter("run", wfBean.getRun()); 144 query.setParameter("status", wfBean.getStatus().toString()); 145 query.setParameter("wfInstance", wfBean.getWfInstanceBlob()); 146 query.setParameter("lastModTime", wfBean.getLastModifiedTimestamp()); 147 query.setParameter("id", wfBean.getId()); 148 break; 149 default: 150 throw new JPAExecutorException(ErrorCode.E0603, "QueryExecutor cannot set parameters for " 151 + namedQuery.name()); 152 } 153 return query; 154 } 155 156 @Override 157 public Query getSelectQuery(WorkflowJobQuery namedQuery, EntityManager em, Object... parameters) 158 throws JPAExecutorException { 159 Query query = em.createNamedQuery(namedQuery.name()); 160 switch (namedQuery) { 161 case GET_WORKFLOW: 162 case GET_WORKFLOW_STARTTIME: 163 case GET_WORKFLOW_START_END_TIME: 164 case GET_WORKFLOW_USER_GROUP: 165 case GET_WORKFLOW_SUSPEND: 166 case GET_WORKFLOW_ACTION_OP: 167 case GET_WORKFLOW_RERUN: 168 case GET_WORKFLOW_DEFINITION: 169 case GET_WORKFLOW_KILL: 170 case GET_WORKFLOW_RESUME: 171 case GET_WORKFLOW_STATUS: 172 query.setParameter("id", parameters[0]); 173 break; 174 case GET_WORKFLOWS_PARENT_COORD_RERUN: 175 query.setParameter("parentId", parameters[0]); 176 break; 177 case GET_COMPLETED_COORD_WORKFLOWS_OLDER_THAN: 178 long dayInMs = 24 * 60 * 60 * 1000; 179 long olderThanDays = (Long) parameters[0]; 180 Timestamp maxEndtime = new Timestamp(System.currentTimeMillis() - (olderThanDays * dayInMs)); 181 query.setParameter("endTime", maxEndtime); 182 query.setFirstResult((Integer) parameters[1]); 183 query.setMaxResults((Integer) parameters[2]); 184 break; 185 default: 186 throw new JPAExecutorException(ErrorCode.E0603, "QueryExecutor cannot set parameters for " 187 + namedQuery.name()); 188 } 189 return query; 190 } 191 192 @Override 193 public int executeUpdate(WorkflowJobQuery namedQuery, WorkflowJobBean jobBean) throws JPAExecutorException { 194 JPAService jpaService = Services.get().get(JPAService.class); 195 EntityManager em = jpaService.getEntityManager(); 196 Query query = getUpdateQuery(namedQuery, jobBean, em); 197 int ret = jpaService.executeUpdate(namedQuery.name(), query, em); 198 return ret; 199 } 200 201 private WorkflowJobBean constructBean(WorkflowJobQuery namedQuery, Object ret, Object... parameters) 202 throws JPAExecutorException { 203 WorkflowJobBean bean; 204 Object[] arr; 205 switch (namedQuery) { 206 case GET_WORKFLOW: 207 bean = (WorkflowJobBean) ret; 208 break; 209 case GET_WORKFLOW_STARTTIME: 210 bean = new WorkflowJobBean(); 211 arr = (Object[]) ret; 212 bean.setId((String) arr[0]); 213 bean.setStartTime(DateUtils.toDate((Timestamp) arr[1])); 214 break; 215 case GET_WORKFLOW_START_END_TIME: 216 bean = new WorkflowJobBean(); 217 arr = (Object[]) ret; 218 bean.setId((String) arr[0]); 219 bean.setStartTime(DateUtils.toDate((Timestamp) arr[1])); 220 bean.setEndTime(DateUtils.toDate((Timestamp) arr[2])); 221 break; 222 case GET_WORKFLOW_USER_GROUP: 223 bean = new WorkflowJobBean(); 224 arr = (Object[]) ret; 225 bean.setUser((String) arr[0]); 226 bean.setGroup((String) arr[1]); 227 break; 228 case GET_WORKFLOW_SUSPEND: 229 bean = new WorkflowJobBean(); 230 arr = (Object[]) ret; 231 bean.setId((String) arr[0]); 232 bean.setUser((String) arr[1]); 233 bean.setGroup((String) arr[2]); 234 bean.setAppName((String) arr[3]); 235 bean.setStatusStr((String) arr[4]); 236 bean.setParentId((String) arr[5]); 237 bean.setStartTime(DateUtils.toDate((Timestamp) arr[6])); 238 bean.setEndTime(DateUtils.toDate((Timestamp) arr[7])); 239 bean.setLogToken((String) arr[8]); 240 bean.setWfInstanceBlob((BinaryBlob) (arr[9])); 241 break; 242 case GET_WORKFLOW_ACTION_OP: 243 bean = new WorkflowJobBean(); 244 arr = (Object[]) ret; 245 bean.setId((String) arr[0]); 246 bean.setUser((String) arr[1]); 247 bean.setGroup((String) arr[2]); 248 bean.setAppName((String) arr[3]); 249 bean.setAppPath((String) arr[4]); 250 bean.setStatusStr((String) arr[5]); 251 bean.setRun((Integer) arr[6]); 252 bean.setParentId((String) arr[7]); 253 bean.setLogToken((String) arr[8]); 254 bean.setWfInstanceBlob((BinaryBlob) (arr[9])); 255 bean.setProtoActionConfBlob((StringBlob) arr[10]); 256 break; 257 case GET_WORKFLOW_RERUN: 258 bean = new WorkflowJobBean(); 259 arr = (Object[]) ret; 260 bean.setId((String) arr[0]); 261 bean.setUser((String) arr[1]); 262 bean.setGroup((String) arr[2]); 263 bean.setAppName((String) arr[3]); 264 bean.setStatusStr((String) arr[4]); 265 bean.setRun((Integer) arr[5]); 266 bean.setLogToken((String) arr[6]); 267 bean.setWfInstanceBlob((BinaryBlob) (arr[7])); 268 bean.setParentId((String)arr[8]); 269 break; 270 case GET_WORKFLOW_DEFINITION: 271 bean = new WorkflowJobBean(); 272 arr = (Object[]) ret; 273 bean.setId((String) arr[0]); 274 bean.setUser((String) arr[1]); 275 bean.setGroup((String) arr[2]); 276 bean.setAppName((String) arr[3]); 277 bean.setLogToken((String) arr[4]); 278 bean.setWfInstanceBlob((BinaryBlob) (arr[5])); 279 break; 280 case GET_WORKFLOW_KILL: 281 bean = new WorkflowJobBean(); 282 arr = (Object[]) ret; 283 bean.setId((String) arr[0]); 284 bean.setUser((String) arr[1]); 285 bean.setGroup((String) arr[2]); 286 bean.setAppName((String) arr[3]); 287 bean.setAppPath((String) arr[4]); 288 bean.setStatusStr((String) arr[5]); 289 bean.setParentId((String) arr[6]); 290 bean.setStartTime(DateUtils.toDate((Timestamp) arr[7])); 291 bean.setEndTime(DateUtils.toDate((Timestamp) arr[8])); 292 bean.setLogToken((String) arr[9]); 293 bean.setWfInstanceBlob((BinaryBlob) (arr[10])); 294 bean.setSlaXmlBlob((StringBlob) arr[11]); 295 bean.setProtoActionConfBlob((StringBlob) arr[12]); 296 break; 297 case GET_WORKFLOW_RESUME: 298 bean = new WorkflowJobBean(); 299 arr = (Object[]) ret; 300 bean.setId((String) arr[0]); 301 bean.setUser((String) arr[1]); 302 bean.setGroup((String) arr[2]); 303 bean.setAppName((String) arr[3]); 304 bean.setAppPath((String) arr[4]); 305 bean.setStatusStr((String) arr[5]); 306 bean.setParentId((String) arr[6]); 307 bean.setStartTime(DateUtils.toDate((Timestamp) arr[7])); 308 bean.setEndTime(DateUtils.toDate((Timestamp) arr[8])); 309 bean.setLogToken((String) arr[9]); 310 bean.setWfInstanceBlob((BinaryBlob) (arr[10])); 311 bean.setProtoActionConfBlob((StringBlob) arr[11]); 312 break; 313 case GET_WORKFLOW_STATUS: 314 bean = new WorkflowJobBean(); 315 bean.setId((String) parameters[0]); 316 bean.setStatusStr((String) ret); 317 break; 318 case GET_WORKFLOWS_PARENT_COORD_RERUN: 319 bean = new WorkflowJobBean(); 320 arr = (Object[]) ret; 321 bean.setId((String) arr[0]); 322 bean.setStatusStr((String) arr[1]); 323 bean.setStartTime(DateUtils.toDate((Timestamp) arr[2])); 324 bean.setEndTime(DateUtils.toDate((Timestamp) arr[3])); 325 break; 326 case GET_COMPLETED_COORD_WORKFLOWS_OLDER_THAN: 327 bean = new WorkflowJobBean(); 328 arr = (Object[]) ret; 329 bean.setId((String) arr[0]); 330 bean.setParentId((String) arr[1]); 331 break; 332 default: 333 throw new JPAExecutorException(ErrorCode.E0603, "QueryExecutor cannot construct job bean for " 334 + namedQuery.name()); 335 } 336 return bean; 337 } 338 339 @Override 340 public WorkflowJobBean get(WorkflowJobQuery namedQuery, Object... parameters) throws JPAExecutorException { 341 JPAService jpaService = Services.get().get(JPAService.class); 342 EntityManager em = jpaService.getEntityManager(); 343 Query query = getSelectQuery(namedQuery, em, parameters); 344 Object ret = jpaService.executeGet(namedQuery.name(), query, em); 345 if (ret == null) { 346 throw new JPAExecutorException(ErrorCode.E0604, query.toString()); 347 } 348 WorkflowJobBean bean = constructBean(namedQuery, ret, parameters); 349 return bean; 350 } 351 352 @Override 353 public List<WorkflowJobBean> getList(WorkflowJobQuery namedQuery, Object... parameters) throws JPAExecutorException { 354 JPAService jpaService = Services.get().get(JPAService.class); 355 EntityManager em = jpaService.getEntityManager(); 356 Query query = getSelectQuery(namedQuery, em, parameters); 357 List<?> retList = (List<?>) jpaService.executeGetList(namedQuery.name(), query, em); 358 List<WorkflowJobBean> beanList = new ArrayList<WorkflowJobBean>(); 359 if (retList != null) { 360 for (Object ret : retList) { 361 beanList.add(constructBean(namedQuery, ret, parameters)); 362 } 363 } 364 return beanList; 365 } 366 367 @Override 368 public Object getSingleValue(WorkflowJobQuery namedQuery, Object... parameters) throws JPAExecutorException { 369 throw new UnsupportedOperationException(); 370 } 371}