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.wf; 020 021import org.apache.hadoop.conf.Configuration; 022import org.apache.hadoop.fs.Path; 023import org.apache.hadoop.fs.FileSystem; 024import org.apache.oozie.AppType; 025import org.apache.oozie.SLAEventBean; 026import org.apache.oozie.command.wf.ActionXCommand.ActionExecutorContext; 027import org.apache.oozie.WorkflowJobBean; 028import org.apache.oozie.ErrorCode; 029import org.apache.oozie.action.oozie.SubWorkflowActionExecutor; 030import org.apache.oozie.service.HadoopAccessorException; 031import org.apache.oozie.service.JPAService; 032import org.apache.oozie.service.UUIDService; 033import org.apache.oozie.service.WorkflowStoreService; 034import org.apache.oozie.service.WorkflowAppService; 035import org.apache.oozie.service.HadoopAccessorService; 036import org.apache.oozie.service.Services; 037import org.apache.oozie.service.DagXLogInfoService; 038import org.apache.oozie.util.ELUtils; 039import org.apache.oozie.util.LogUtils; 040import org.apache.oozie.sla.SLAOperations; 041import org.apache.oozie.util.XLog; 042import org.apache.oozie.util.ParamChecker; 043import org.apache.oozie.util.XConfiguration; 044import org.apache.oozie.util.XmlUtils; 045import org.apache.oozie.command.CommandException; 046import org.apache.oozie.executor.jpa.BatchQueryExecutor; 047import org.apache.oozie.executor.jpa.JPAExecutorException; 048import org.apache.oozie.service.ELService; 049import org.apache.oozie.store.StoreException; 050import org.apache.oozie.workflow.WorkflowApp; 051import org.apache.oozie.workflow.WorkflowException; 052import org.apache.oozie.workflow.WorkflowInstance; 053import org.apache.oozie.workflow.WorkflowLib; 054import org.apache.oozie.util.ELEvaluator; 055import org.apache.oozie.util.InstrumentUtils; 056import org.apache.oozie.util.PropertiesUtils; 057import org.apache.oozie.util.db.SLADbOperations; 058import org.apache.oozie.service.SchemaService.SchemaName; 059import org.apache.oozie.client.OozieClient; 060import org.apache.oozie.client.WorkflowJob; 061import org.apache.oozie.client.SLAEvent.SlaAppType; 062import org.apache.oozie.client.rest.JsonBean; 063import org.jdom.Element; 064import org.jdom.filter.ElementFilter; 065 066import java.util.ArrayList; 067import java.util.Date; 068import java.util.Iterator; 069import java.util.List; 070import java.util.Map; 071import java.util.Set; 072import java.util.HashSet; 073import java.io.IOException; 074import java.net.URI; 075 076@SuppressWarnings("deprecation") 077public class SubmitXCommand extends WorkflowXCommand<String> { 078 public static final String CONFIG_DEFAULT = "config-default.xml"; 079 080 private Configuration conf; 081 private List<JsonBean> insertList = new ArrayList<JsonBean>(); 082 private String parentId; 083 084 /** 085 * Constructor to create the workflow Submit Command. 086 * 087 * @param conf : Configuration for workflow job 088 */ 089 public SubmitXCommand(Configuration conf) { 090 super("submit", "submit", 1); 091 this.conf = ParamChecker.notNull(conf, "conf"); 092 } 093 094 /** 095 * Constructor for submitting wf through coordinator 096 * 097 * @param conf : Configuration for workflow job 098 * @param parentId: the coord action id 099 */ 100 public SubmitXCommand(Configuration conf, String parentId) { 101 this(conf); 102 this.parentId = parentId; 103 } 104 105 /** 106 * Constructor to create the workflow Submit Command. 107 * 108 * @param dryrun : if dryrun 109 * @param conf : Configuration for workflow job 110 * @param authToken : To be used for authentication 111 */ 112 public SubmitXCommand(boolean dryrun, Configuration conf) { 113 this(conf); 114 this.dryrun = dryrun; 115 } 116 117 private static final Set<String> DISALLOWED_USER_PROPERTIES = new HashSet<String>(); 118 119 static { 120 String[] badUserProps = {PropertiesUtils.DAYS, PropertiesUtils.HOURS, PropertiesUtils.MINUTES, 121 PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, PropertiesUtils.TB, PropertiesUtils.PB, 122 PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN, 123 PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS}; 124 PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_USER_PROPERTIES); 125 } 126 127 @Override 128 protected String execute() throws CommandException { 129 InstrumentUtils.incrJobCounter(getName(), 1, getInstrumentation()); 130 WorkflowAppService wps = Services.get().get(WorkflowAppService.class); 131 try { 132 XLog.Info.get().setParameter(DagXLogInfoService.TOKEN, conf.get(OozieClient.LOG_TOKEN)); 133 String user = conf.get(OozieClient.USER_NAME); 134 URI uri = new URI(conf.get(OozieClient.APP_PATH)); 135 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 136 Configuration fsConf = has.createJobConf(uri.getAuthority()); 137 FileSystem fs = has.createFileSystem(user, uri, fsConf); 138 139 Path configDefault = null; 140 Configuration defaultConf = null; 141 // app path could be a directory 142 Path path = new Path(uri.getPath()); 143 if (!fs.isFile(path)) { 144 configDefault = new Path(path, CONFIG_DEFAULT); 145 } else { 146 configDefault = new Path(path.getParent(), CONFIG_DEFAULT); 147 } 148 149 if (fs.exists(configDefault)) { 150 try { 151 defaultConf = new XConfiguration(fs.open(configDefault)); 152 PropertiesUtils.checkDisallowedProperties(defaultConf, DISALLOWED_USER_PROPERTIES); 153 PropertiesUtils.checkDefaultDisallowedProperties(defaultConf); 154 XConfiguration.injectDefaults(defaultConf, conf); 155 } 156 catch (IOException ex) { 157 throw new IOException("default configuration file, " + ex.getMessage(), ex); 158 } 159 } 160 161 WorkflowApp app = wps.parseDef(conf, defaultConf); 162 XConfiguration protoActionConf = wps.createProtoActionConf(conf, true); 163 WorkflowLib workflowLib = Services.get().get(WorkflowStoreService.class).getWorkflowLibWithNoDB(); 164 165 PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_USER_PROPERTIES); 166 167 // Resolving all variables in the job properties. 168 // This ensures the Hadoop Configuration semantics is preserved. 169 XConfiguration resolvedVarsConf = new XConfiguration(); 170 for (Map.Entry<String, String> entry : conf) { 171 resolvedVarsConf.set(entry.getKey(), conf.get(entry.getKey())); 172 } 173 conf = resolvedVarsConf; 174 175 WorkflowInstance wfInstance; 176 try { 177 wfInstance = workflowLib.createInstance(app, conf); 178 } 179 catch (WorkflowException e) { 180 throw new StoreException(e); 181 } 182 183 Configuration conf = wfInstance.getConf(); 184 // System.out.println("WF INSTANCE CONF:"); 185 // System.out.println(XmlUtils.prettyPrint(conf).toString()); 186 187 WorkflowJobBean workflow = new WorkflowJobBean(); 188 workflow.setId(wfInstance.getId()); 189 workflow.setAppName(ELUtils.resolveAppName(app.getName(), conf)); 190 workflow.setAppPath(conf.get(OozieClient.APP_PATH)); 191 workflow.setConf(XmlUtils.prettyPrint(conf).toString()); 192 workflow.setProtoActionConf(protoActionConf.toXmlString()); 193 workflow.setCreatedTime(new Date()); 194 workflow.setLastModifiedTime(new Date()); 195 workflow.setLogToken(conf.get(OozieClient.LOG_TOKEN, "")); 196 workflow.setStatus(WorkflowJob.Status.PREP); 197 workflow.setRun(0); 198 workflow.setUser(conf.get(OozieClient.USER_NAME)); 199 workflow.setGroup(conf.get(OozieClient.GROUP_NAME)); 200 workflow.setWorkflowInstance(wfInstance); 201 workflow.setExternalId(conf.get(OozieClient.EXTERNAL_ID)); 202 // Set parent id if it doesn't already have one (for subworkflows) 203 if (workflow.getParentId() == null) { 204 workflow.setParentId(conf.get(SubWorkflowActionExecutor.PARENT_ID)); 205 } 206 // Set to coord action Id if workflow submitted through coordinator 207 if (workflow.getParentId() == null) { 208 workflow.setParentId(parentId); 209 } 210 211 LogUtils.setLogInfo(workflow); 212 LOG.debug("Workflow record created, Status [{0}]", workflow.getStatus()); 213 Element wfElem = XmlUtils.parseXml(app.getDefinition()); 214 ELEvaluator evalSla = createELEvaluatorForGroup(conf, "wf-sla-submit"); 215 String jobSlaXml = verifySlaElements(wfElem, evalSla); 216 if (!dryrun) { 217 writeSLARegistration(wfElem, jobSlaXml, workflow.getId(), workflow.getParentId(), workflow.getUser(), 218 workflow.getGroup(), workflow.getAppName(), LOG, evalSla); 219 workflow.setSlaXml(jobSlaXml); 220 // System.out.println("SlaXml :"+ slaXml); 221 222 //store.insertWorkflow(workflow); 223 insertList.add(workflow); 224 JPAService jpaService = Services.get().get(JPAService.class); 225 if (jpaService != null) { 226 try { 227 BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, null, null); 228 } 229 catch (JPAExecutorException je) { 230 throw new CommandException(je); 231 } 232 } 233 else { 234 LOG.error(ErrorCode.E0610); 235 return null; 236 } 237 238 return workflow.getId(); 239 } 240 else { 241 // Checking variable substitution for dryrun 242 ActionExecutorContext context = new ActionXCommand.ActionExecutorContext(workflow, null, false, false); 243 Element workflowXml = XmlUtils.parseXml(app.getDefinition()); 244 removeSlaElements(workflowXml); 245 String workflowXmlString = XmlUtils.removeComments(XmlUtils.prettyPrint(workflowXml).toString()); 246 workflowXmlString = context.getELEvaluator().evaluate(workflowXmlString, String.class); 247 workflowXml = XmlUtils.parseXml(workflowXmlString); 248 249 Iterator<Element> it = workflowXml.getDescendants(new ElementFilter("job-xml")); 250 251 // Checking all variable substitutions in job-xml files 252 while (it.hasNext()) { 253 Element e = it.next(); 254 String jobXml = e.getTextTrim(); 255 Path xmlPath = new Path(workflow.getAppPath(), jobXml); 256 Configuration jobXmlConf = new XConfiguration(fs.open(xmlPath)); 257 258 259 String jobXmlConfString = XmlUtils.prettyPrint(jobXmlConf).toString(); 260 jobXmlConfString = XmlUtils.removeComments(jobXmlConfString); 261 context.getELEvaluator().evaluate(jobXmlConfString, String.class); 262 } 263 264 return "OK"; 265 } 266 } 267 catch (WorkflowException ex) { 268 throw new CommandException(ex); 269 } 270 catch (HadoopAccessorException ex) { 271 throw new CommandException(ex); 272 } 273 catch (Exception ex) { 274 throw new CommandException(ErrorCode.E0803, ex.getMessage(), ex); 275 } 276 } 277 278 private void removeSlaElements(Element eWfJob) { 279 Element sla = XmlUtils.getSLAElement(eWfJob); 280 if (sla != null) { 281 eWfJob.removeChildren(sla.getName(), sla.getNamespace()); 282 } 283 284 for (Element action : (List<Element>) eWfJob.getChildren("action", eWfJob.getNamespace())) { 285 sla = XmlUtils.getSLAElement(action); 286 if (sla != null) { 287 action.removeChildren(sla.getName(), sla.getNamespace()); 288 } 289 } 290 } 291 private String verifySlaElements(Element eWfJob, ELEvaluator evalSla) throws CommandException { 292 String jobSlaXml = ""; 293 // Validate WF job 294 Element eSla = XmlUtils.getSLAElement(eWfJob); 295 if (eSla != null) { 296 jobSlaXml = resolveSla(eSla, evalSla); 297 } 298 299 // Validate all actions 300 for (Element action : (List<Element>) eWfJob.getChildren("action", eWfJob.getNamespace())) { 301 eSla = XmlUtils.getSLAElement(action); 302 if (eSla != null) { 303 resolveSla(eSla, evalSla); 304 } 305 } 306 return jobSlaXml; 307 } 308 309 private void writeSLARegistration(Element eWfJob, String slaXml, String jobId, String parentId, String user, 310 String group, String appName, XLog log, ELEvaluator evalSla) throws CommandException { 311 try { 312 if (slaXml != null && slaXml.length() > 0) { 313 Element eSla = XmlUtils.parseXml(slaXml); 314 SLAEventBean slaEvent = SLADbOperations.createSlaRegistrationEvent(eSla, jobId, 315 SlaAppType.WORKFLOW_JOB, user, group, log); 316 if(slaEvent != null) { 317 insertList.add(slaEvent); 318 } 319 // insert into new table 320 SLAOperations.createSlaRegistrationEvent(eSla, jobId, parentId, AppType.WORKFLOW_JOB, user, appName, 321 log, false); 322 } 323 // Add sla for wf actions 324 for (Element action : (List<Element>) eWfJob.getChildren("action", eWfJob.getNamespace())) { 325 Element actionSla = XmlUtils.getSLAElement(action); 326 if (actionSla != null) { 327 String actionSlaXml = SubmitXCommand.resolveSla(actionSla, evalSla); 328 actionSla = XmlUtils.parseXml(actionSlaXml); 329 String actionId = Services.get().get(UUIDService.class) 330 .generateChildId(jobId, action.getAttributeValue("name") + ""); 331 SLAOperations.createSlaRegistrationEvent(actionSla, actionId, jobId, AppType.WORKFLOW_ACTION, 332 user, appName, log, false); 333 } 334 } 335 } 336 catch (Exception e) { 337 e.printStackTrace(); 338 throw new CommandException(ErrorCode.E1007, "workflow " + jobId, e.getMessage(), e); 339 } 340 } 341 342 /** 343 * Resolve variables in sla xml element. 344 * 345 * @param eSla sla xml element 346 * @param evalSla sla evaluator 347 * @return sla xml string after evaluation 348 * @throws CommandException 349 */ 350 public static String resolveSla(Element eSla, ELEvaluator evalSla) throws CommandException { 351 // EL evaluation 352 String slaXml = XmlUtils.prettyPrint(eSla).toString(); 353 try { 354 slaXml = XmlUtils.removeComments(slaXml); 355 slaXml = evalSla.evaluate(slaXml, String.class); 356 XmlUtils.validateData(slaXml, SchemaName.SLA_ORIGINAL); 357 return slaXml; 358 } 359 catch (Exception e) { 360 throw new CommandException(ErrorCode.E1004, "Validation error :" + e.getMessage(), e); 361 } 362 } 363 364 /** 365 * Create an EL evaluator for a given group. 366 * 367 * @param conf configuration variable 368 * @param group group variable 369 * @return the evaluator created for the group 370 */ 371 public static ELEvaluator createELEvaluatorForGroup(Configuration conf, String group) { 372 ELEvaluator eval = Services.get().get(ELService.class).createEvaluator(group); 373 for (Map.Entry<String, String> entry : conf) { 374 eval.setVariable(entry.getKey(), entry.getValue()); 375 } 376 return eval; 377 } 378 379 @Override 380 public String getEntityKey() { 381 return null; 382 } 383 384 @Override 385 protected boolean isLockRequired() { 386 return false; 387 } 388 389 @Override 390 protected void loadState() { 391 392 } 393 394 @Override 395 protected void verifyPrecondition() throws CommandException { 396 397 } 398 399}