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.command.coord; 020 021import org.apache.hadoop.conf.Configuration; 022import org.apache.oozie.action.hadoop.OozieJobInfo; 023import org.apache.oozie.client.CoordinatorAction; 024import org.apache.oozie.client.OozieClient; 025import org.apache.oozie.CoordinatorActionBean; 026import org.apache.oozie.DagEngineException; 027import org.apache.oozie.DagEngine; 028import org.apache.oozie.ErrorCode; 029import org.apache.oozie.SLAEventBean; 030import org.apache.oozie.WorkflowJobBean; 031import org.apache.oozie.command.CommandException; 032import org.apache.oozie.command.PreconditionException; 033import org.apache.oozie.service.DagEngineService; 034import org.apache.oozie.service.EventHandlerService; 035import org.apache.oozie.service.JPAService; 036import org.apache.oozie.service.Services; 037import org.apache.oozie.util.ConfigUtils; 038import org.apache.oozie.util.*; 039import org.apache.oozie.util.db.SLADbOperations; 040import org.apache.oozie.client.SLAEvent.SlaAppType; 041import org.apache.oozie.client.SLAEvent.Status; 042import org.apache.oozie.client.rest.JsonBean; 043import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry; 044import org.apache.oozie.executor.jpa.BatchQueryExecutor; 045import org.apache.oozie.executor.jpa.CoordActionQueryExecutor.CoordActionQuery; 046import org.apache.oozie.executor.jpa.JPAExecutorException; 047import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor; 048import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor.WorkflowJobQuery; 049import org.jdom.Element; 050import org.jdom.JDOMException; 051 052import java.io.IOException; 053import java.io.StringReader; 054import java.util.ArrayList; 055import java.util.Date; 056import java.util.List; 057 058@SuppressWarnings("deprecation") 059public class CoordActionStartXCommand extends CoordinatorXCommand<Void> { 060 061 public static final String EL_ERROR = "EL_ERROR"; 062 public static final String EL_EVAL_ERROR = "EL_EVAL_ERROR"; 063 public static final String COULD_NOT_START = "COULD_NOT_START"; 064 public static final String START_DATA_MISSING = "START_DATA_MISSING"; 065 public static final String EXEC_DATA_MISSING = "EXEC_DATA_MISSING"; 066 public static final String OOZIE_COORD_ACTION_NOMINAL_TIME = "oozie.coord.action.nominal_time"; 067 068 private final XLog log = getLog(); 069 private String actionId = null; 070 private String user = null; 071 private String appName = null; 072 private CoordinatorActionBean coordAction = null; 073 private JPAService jpaService = null; 074 private String jobId = null; 075 private List<UpdateEntry> updateList = new ArrayList<UpdateEntry>(); 076 private List<JsonBean> insertList = new ArrayList<JsonBean>(); 077 078 public CoordActionStartXCommand(String id, String user, String appName, String jobId) { 079 //super("coord_action_start", "coord_action_start", 1, XLog.OPS); 080 super("coord_action_start", "coord_action_start", 1); 081 this.actionId = ParamChecker.notEmpty(id, "id"); 082 this.user = ParamChecker.notEmpty(user, "user"); 083 this.appName = ParamChecker.notEmpty(appName, "appName"); 084 this.jobId = jobId; 085 } 086 087 @Override 088 protected void setLogInfo() { 089 LogUtils.setLogInfo(actionId); 090 } 091 092 /** 093 * Create config to pass to WF Engine 1. Get createdConf from coord_actions table 2. Get actionXml from 094 * coord_actions table. Extract all 'property' tags and merge createdConf (overwrite duplicate keys). 3. Extract 095 * 'app-path' from actionXML. Create a new property called 'oozie.wf.application.path' and merge with createdConf 096 * (overwrite duplicate keys) 4. Read contents of config-default.xml in workflow directory. 5. Merge createdConf 097 * with config-default.xml (overwrite duplicate keys). 6. Results is runConf which is saved in coord_actions table. 098 * Merge Action createdConf with actionXml to create new runConf with replaced variables 099 * 100 * @param action CoordinatorActionBean 101 * @return Configuration 102 * @throws CommandException 103 */ 104 private Configuration mergeConfig(CoordinatorActionBean action) throws CommandException { 105 String createdConf = action.getCreatedConf(); 106 String actionXml = action.getActionXml(); 107 Element workflowProperties = null; 108 try { 109 workflowProperties = XmlUtils.parseXml(actionXml); 110 } 111 catch (JDOMException e1) { 112 log.warn("Configuration parse error in:" + actionXml); 113 throw new CommandException(ErrorCode.E1005, e1.getMessage(), e1); 114 } 115 // generate the 'runConf' for this action 116 // Step 1: runConf = createdConf 117 Configuration runConf = null; 118 try { 119 runConf = new XConfiguration(new StringReader(createdConf)); 120 } 121 catch (IOException e1) { 122 log.warn("Configuration parse error in:" + createdConf); 123 throw new CommandException(ErrorCode.E1005, e1.getMessage(), e1); 124 } 125 // Step 2: Merge local properties into runConf 126 // extract 'property' tags under 'configuration' block in the 127 // coordinator.xml (saved in actionxml column) 128 // convert Element to XConfiguration 129 Element configElement = workflowProperties.getChild("action", workflowProperties.getNamespace()) 130 .getChild("workflow", workflowProperties.getNamespace()).getChild("configuration", 131 workflowProperties.getNamespace()); 132 if (configElement != null) { 133 String strConfig = XmlUtils.prettyPrint(configElement).toString(); 134 Configuration localConf; 135 try { 136 localConf = new XConfiguration(new StringReader(strConfig)); 137 } 138 catch (IOException e1) { 139 log.warn("Configuration parse error in:" + strConfig); 140 throw new CommandException(ErrorCode.E1005, e1.getMessage(), e1); 141 } 142 143 // copy configuration properties in coordinator.xml to the runConf 144 XConfiguration.copy(localConf, runConf); 145 } 146 147 // Step 3: Extract value of 'app-path' in actionxml, and save it as a 148 // new property called 'oozie.wf.application.path' 149 // WF Engine requires the path to the workflow.xml to be saved under 150 // this property name 151 String appPath = workflowProperties.getChild("action", workflowProperties.getNamespace()).getChild("workflow", 152 workflowProperties.getNamespace()).getChild("app-path", workflowProperties.getNamespace()).getValue(); 153 runConf.set("oozie.wf.application.path", appPath); 154 155 ConfigUtils.checkAndSetDisallowedProperties(runConf, 156 this.user, 157 new CommandException(ErrorCode.E1003, 158 String.format("%s=%s", OozieClient.USER_NAME, runConf.get(OozieClient.USER_NAME))), 159 true); 160 161 return runConf; 162 } 163 164 @Override 165 protected Void execute() throws CommandException { 166 boolean makeFail = true; 167 String errCode = ""; 168 String errMsg = ""; 169 ParamChecker.notEmpty(user, "user"); 170 171 log.debug("actionid=" + actionId + ", status=" + coordAction.getStatus()); 172 if (coordAction.getStatus() == CoordinatorAction.Status.SUBMITTED) { 173 // log.debug("getting.. job id: " + coordAction.getJobId()); 174 // create merged runConf to pass to WF Engine 175 Configuration runConf = mergeConfig(coordAction); 176 coordAction.setRunConf(XmlUtils.prettyPrint(runConf).toString()); 177 // log.debug("%%% merged runconf=" + 178 // XmlUtils.prettyPrint(runConf).toString()); 179 DagEngine dagEngine = Services.get().get(DagEngineService.class).getDagEngine(user); 180 try { 181 Configuration conf = new XConfiguration(new StringReader(coordAction.getRunConf())); 182 SLAEventBean slaEvent = SLADbOperations.createStatusEvent(coordAction.getSlaXml(), coordAction.getId(), Status.STARTED, 183 SlaAppType.COORDINATOR_ACTION, log); 184 if(slaEvent != null) { 185 insertList.add(slaEvent); 186 } 187 if (OozieJobInfo.isJobInfoEnabled()) { 188 conf.set(OozieJobInfo.COORD_ID, actionId); 189 conf.set(OozieJobInfo.COORD_NAME, appName); 190 conf.set(OozieJobInfo.COORD_NOMINAL_TIME, coordAction.getNominalTimestamp().toString()); 191 } 192 // Normalize workflow appPath here; 193 JobUtils.normalizeAppPath(conf.get(OozieClient.USER_NAME), conf.get(OozieClient.GROUP_NAME), conf); 194 if (coordAction.getExternalId() != null) { 195 conf.setBoolean(OozieClient.RERUN_FAIL_NODES, true); 196 dagEngine.reRun(coordAction.getExternalId(), conf); 197 } else { 198 // Pushing the nominal time in conf to use for launcher tag search 199 conf.set(OOZIE_COORD_ACTION_NOMINAL_TIME,String.valueOf(coordAction.getNominalTime().getTime())); 200 String wfId = dagEngine.submitJobFromCoordinator(conf, actionId); 201 coordAction.setExternalId(wfId); 202 } 203 coordAction.setStatus(CoordinatorAction.Status.RUNNING); 204 coordAction.incrementAndGetPending(); 205 206 //store.updateCoordinatorAction(coordAction); 207 JPAService jpaService = Services.get().get(JPAService.class); 208 if (jpaService != null) { 209 log.debug("Updating WF record for WFID :" + coordAction.getExternalId() + " with parent id: " + actionId); 210 WorkflowJobBean wfJob = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW_STARTTIME, coordAction.getExternalId()); 211 wfJob.setParentId(actionId); 212 wfJob.setLastModifiedTime(new Date()); 213 BatchQueryExecutor executor = BatchQueryExecutor.getInstance(); 214 updateList.add(new UpdateEntry<WorkflowJobQuery>( 215 WorkflowJobQuery.UPDATE_WORKFLOW_PARENT_MODIFIED, wfJob)); 216 updateList.add(new UpdateEntry<CoordActionQuery>( 217 CoordActionQuery.UPDATE_COORD_ACTION_FOR_START, coordAction)); 218 try { 219 executor.executeBatchInsertUpdateDelete(insertList, updateList, null); 220 queue(new CoordActionNotificationXCommand(coordAction), 100); 221 if (EventHandlerService.isEnabled()) { 222 generateEvent(coordAction, user, appName, wfJob.getStartTime()); 223 } 224 } 225 catch (JPAExecutorException je) { 226 throw new CommandException(je); 227 } 228 } 229 else { 230 log.error(ErrorCode.E0610); 231 } 232 233 makeFail = false; 234 } 235 catch (DagEngineException dee) { 236 errMsg = dee.getMessage(); 237 errCode = dee.getErrorCode().toString(); 238 log.warn("can not create DagEngine for submitting jobs", dee); 239 } 240 catch (CommandException ce) { 241 errMsg = ce.getMessage(); 242 errCode = ce.getErrorCode().toString(); 243 log.warn("command exception occured ", ce); 244 } 245 catch (java.io.IOException ioe) { 246 errMsg = ioe.getMessage(); 247 errCode = "E1005"; 248 log.warn("Configuration parse error. read from DB :" + coordAction.getRunConf(), ioe); 249 } 250 catch (Exception ex) { 251 errMsg = ex.getMessage(); 252 errCode = "E1005"; 253 log.warn("can not create DagEngine for submitting jobs", ex); 254 } 255 finally { 256 if (makeFail == true) { // No DB exception occurs 257 log.error("Failing the action " + coordAction.getId() + ". Because " + errCode + " : " + errMsg); 258 coordAction.setStatus(CoordinatorAction.Status.FAILED); 259 if (errMsg.length() > 254) { // Because table column size is 255 260 errMsg = errMsg.substring(0, 255); 261 } 262 coordAction.setErrorMessage(errMsg); 263 coordAction.setErrorCode(errCode); 264 265 updateList = new ArrayList<UpdateEntry>(); 266 updateList.add(new UpdateEntry<CoordActionQuery>( 267 CoordActionQuery.UPDATE_COORD_ACTION_FOR_START, coordAction)); 268 insertList = new ArrayList<JsonBean>(); 269 270 SLAEventBean slaEvent = SLADbOperations.createStatusEvent(coordAction.getSlaXml(), coordAction.getId(), Status.FAILED, 271 SlaAppType.COORDINATOR_ACTION, log); 272 if(slaEvent != null) { 273 insertList.add(slaEvent); //Update SLA events 274 } 275 try { 276 // call JPAExecutor to do the bulk writes 277 BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, updateList, null); 278 if (EventHandlerService.isEnabled()) { 279 generateEvent(coordAction, user, appName, null); 280 } 281 } 282 catch (JPAExecutorException je) { 283 throw new CommandException(je); 284 } 285 queue(new CoordActionReadyXCommand(coordAction.getJobId())); 286 } 287 } 288 } 289 return null; 290 } 291 292 @Override 293 public String getEntityKey() { 294 return this.jobId; 295 } 296 297 @Override 298 protected boolean isLockRequired() { 299 return true; 300 } 301 302 @Override 303 protected void loadState() throws CommandException { 304 jpaService = Services.get().get(JPAService.class); 305 try { 306 coordAction = jpaService.execute(new org.apache.oozie.executor.jpa.CoordActionGetForStartJPAExecutor( 307 actionId)); 308 } 309 catch (JPAExecutorException je) { 310 throw new CommandException(je); 311 } 312 LogUtils.setLogInfo(coordAction); 313 } 314 315 @Override 316 protected void verifyPrecondition() throws PreconditionException { 317 if (coordAction.getStatus() != CoordinatorAction.Status.SUBMITTED) { 318 throw new PreconditionException(ErrorCode.E1100, "The coord action [" + actionId + "] must have status " 319 + CoordinatorAction.Status.SUBMITTED.name() + " but has status [" + coordAction.getStatus().name() + "]"); 320 } 321 } 322 323 @Override 324 public String getKey(){ 325 return getName() + "_" + actionId; 326 } 327}