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.oozie.WorkflowJobBean; 023import org.apache.oozie.ErrorCode; 024import org.apache.oozie.service.JPAService; 025import org.apache.oozie.service.WorkflowStoreService; 026import org.apache.oozie.service.WorkflowAppService; 027import org.apache.oozie.service.Services; 028import org.apache.oozie.service.DagXLogInfoService; 029import org.apache.oozie.util.InstrumentUtils; 030import org.apache.oozie.util.LogUtils; 031import org.apache.oozie.util.XLog; 032import org.apache.oozie.util.ParamChecker; 033import org.apache.oozie.util.XConfiguration; 034import org.apache.oozie.util.XmlUtils; 035import org.apache.oozie.command.CommandException; 036import org.apache.oozie.executor.jpa.WorkflowJobInsertJPAExecutor; 037import org.apache.oozie.store.StoreException; 038import org.apache.oozie.workflow.WorkflowApp; 039import org.apache.oozie.workflow.WorkflowException; 040import org.apache.oozie.workflow.WorkflowInstance; 041import org.apache.oozie.workflow.WorkflowLib; 042import org.apache.oozie.util.PropertiesUtils; 043import org.apache.oozie.client.OozieClient; 044import org.apache.oozie.client.WorkflowJob; 045import org.apache.oozie.client.XOozieClient; 046import org.jdom.Element; 047import org.jdom.Namespace; 048 049import java.util.Date; 050import java.util.List; 051import java.util.Map; 052import java.util.Set; 053import java.util.HashSet; 054 055public abstract class SubmitHttpXCommand extends WorkflowXCommand<String> { 056 057 protected static final Set<String> MANDATORY_OOZIE_CONFS = new HashSet<String>(); 058 protected static final Set<String> OPTIONAL_OOZIE_CONFS = new HashSet<String>(); 059 060 static { 061 MANDATORY_OOZIE_CONFS.add(XOozieClient.JT); 062 MANDATORY_OOZIE_CONFS.add(XOozieClient.NN); 063 MANDATORY_OOZIE_CONFS.add(OozieClient.LIBPATH); 064 065 OPTIONAL_OOZIE_CONFS.add(XOozieClient.FILES); 066 OPTIONAL_OOZIE_CONFS.add(XOozieClient.ARCHIVES); 067 } 068 069 private Configuration conf; 070 071 public SubmitHttpXCommand(String name, String type, Configuration conf) { 072 super(name, type, 1); 073 this.conf = ParamChecker.notNull(conf, "conf"); 074 } 075 076 private static final Set<String> DISALLOWED_USER_PROPERTIES = new HashSet<String>(); 077 078 static { 079 String[] badUserProps = { PropertiesUtils.DAYS, PropertiesUtils.HOURS, PropertiesUtils.MINUTES, 080 PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, PropertiesUtils.TB, PropertiesUtils.PB, 081 PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN, 082 PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS }; 083 PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_USER_PROPERTIES); 084 } 085 086 abstract protected Element generateSection(Configuration conf, Namespace ns); 087 088 abstract protected Namespace getSectionNamespace(); 089 090 abstract protected String getWorkflowName(); 091 092 protected void checkMandatoryConf(Configuration conf) { 093 for (String key : MANDATORY_OOZIE_CONFS) { 094 String value = conf.get(key); 095 if (value == null) { 096 throw new RuntimeException(key + " is not specified"); 097 } 098 } 099 } 100 101 protected Namespace getWorkflowNamespace() { 102 return Namespace.getNamespace("uri:oozie:workflow:0.2"); 103 } 104 /** 105 * Generate workflow xml from conf object 106 * 107 * @param conf the configuration object 108 * @return workflow xml def string representation 109 */ 110 protected String getWorkflowXml(Configuration conf) { 111 checkMandatoryConf(conf); 112 113 Namespace ns = getWorkflowNamespace(); 114 Element root = new Element("workflow-app", ns); 115 String name = getWorkflowName(); 116 root.setAttribute("name", "oozie-" + name); 117 118 Element start = new Element("start", ns); 119 String nodeName = name + "1"; 120 start.setAttribute("to", nodeName); 121 root.addContent(start); 122 123 Element action = new Element("action", ns); 124 action.setAttribute("name", nodeName); 125 126 Element ele = generateSection(conf, getSectionNamespace()); 127 action.addContent(ele); 128 129 Element ok = new Element("ok", ns); 130 ok.setAttribute("to", "end"); 131 action.addContent(ok); 132 133 Element error = new Element("error", ns); 134 error.setAttribute("to", "fail"); 135 action.addContent(error); 136 137 root.addContent(action); 138 139 Element kill = new Element("kill", ns); 140 kill.setAttribute("name", "fail"); 141 Element message = new Element("message", ns); 142 message.addContent(name + " failed, error message[${wf:errorMessage(wf:lastErrorNode())}]"); 143 kill.addContent(message); 144 root.addContent(kill); 145 146 Element end = new Element("end", ns); 147 end.setAttribute("name", "end"); 148 root.addContent(end); 149 150 return XmlUtils.prettyPrint(root).toString(); 151 }; 152 153 protected Element generateConfigurationSection(List<String> Dargs, Namespace ns) { 154 Element configuration = new Element("configuration", ns); 155 for (String arg : Dargs) { 156 String name = null, value = null; 157 int pos = arg.indexOf("="); 158 if (pos == -1) { // "-D<name>" or "-D" only 159 name = arg.substring(2, arg.length()); 160 value = ""; 161 } 162 else { // "-D<name>=<value>" 163 name = arg.substring(2, pos); 164 value = arg.substring(pos + 1, arg.length()); 165 } 166 167 Element property = new Element("property", ns); 168 Element nameElement = new Element("name", ns); 169 nameElement.addContent(name); 170 property.addContent(nameElement); 171 Element valueElement = new Element("value", ns); 172 valueElement.addContent(value); 173 property.addContent(valueElement); 174 configuration.addContent(property); 175 } 176 177 return configuration; 178 } 179 180 /* (non-Javadoc) 181 * @see org.apache.oozie.command.XCommand#execute() 182 */ 183 @Override 184 protected String execute() throws CommandException { 185 InstrumentUtils.incrJobCounter(getName(), 1, getInstrumentation()); 186 WorkflowAppService wps = Services.get().get(WorkflowAppService.class); 187 try { 188 XLog.Info.get().setParameter(DagXLogInfoService.TOKEN, conf.get(OozieClient.LOG_TOKEN)); 189 String wfXml = getWorkflowXml(conf); 190 LOG.debug("workflow xml created on the server side is :\n"); 191 LOG.debug(wfXml); 192 WorkflowApp app = wps.parseDef(wfXml, conf); 193 XConfiguration protoActionConf = wps.createProtoActionConf(conf, false); 194 WorkflowLib workflowLib = Services.get().get(WorkflowStoreService.class).getWorkflowLibWithNoDB(); 195 196 PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_USER_PROPERTIES); 197 PropertiesUtils.checkDefaultDisallowedProperties(conf); 198 199 // Resolving all variables in the job properties. 200 // This ensures the Hadoop Configuration semantics is preserved. 201 XConfiguration resolvedVarsConf = new XConfiguration(); 202 for (Map.Entry<String, String> entry : conf) { 203 resolvedVarsConf.set(entry.getKey(), conf.get(entry.getKey())); 204 } 205 conf = resolvedVarsConf; 206 207 WorkflowInstance wfInstance; 208 try { 209 wfInstance = workflowLib.createInstance(app, conf); 210 } 211 catch (WorkflowException e) { 212 throw new StoreException(e); 213 } 214 215 Configuration conf = wfInstance.getConf(); 216 217 WorkflowJobBean workflow = new WorkflowJobBean(); 218 workflow.setId(wfInstance.getId()); 219 workflow.setAppName(app.getName()); 220 workflow.setAppPath(conf.get(OozieClient.APP_PATH)); 221 workflow.setConf(XmlUtils.prettyPrint(conf).toString()); 222 workflow.setProtoActionConf(protoActionConf.toXmlString()); 223 workflow.setCreatedTime(new Date()); 224 workflow.setLastModifiedTime(new Date()); 225 workflow.setLogToken(conf.get(OozieClient.LOG_TOKEN, "")); 226 workflow.setStatus(WorkflowJob.Status.PREP); 227 workflow.setRun(0); 228 workflow.setUser(conf.get(OozieClient.USER_NAME)); 229 workflow.setGroup(conf.get(OozieClient.GROUP_NAME)); 230 workflow.setWorkflowInstance(wfInstance); 231 workflow.setExternalId(conf.get(OozieClient.EXTERNAL_ID)); 232 233 LogUtils.setLogInfo(workflow); 234 JPAService jpaService = Services.get().get(JPAService.class); 235 if (jpaService != null) { 236 jpaService.execute(new WorkflowJobInsertJPAExecutor(workflow)); 237 } 238 else { 239 LOG.error(ErrorCode.E0610); 240 return null; 241 } 242 243 return workflow.getId(); 244 } 245 catch (WorkflowException ex) { 246 throw new CommandException(ex); 247 } 248 catch (Exception ex) { 249 throw new CommandException(ErrorCode.E0803, ex.getMessage(), ex); 250 } 251 } 252 253 static private void addSection(Element X, Namespace ns, String filesStr, String tagName) { 254 if (filesStr != null) { 255 String[] files = filesStr.split(","); 256 for (String f : files) { 257 Element tagElement = new Element(tagName, ns); 258 if (f.contains("#")) { 259 tagElement.addContent(f); 260 } 261 else { 262 String filename = f.substring(f.lastIndexOf("/") + 1, f.length()); 263 if (filename == null || filename.isEmpty()) { 264 tagElement.addContent(f); 265 } 266 else { 267 tagElement.addContent(f + "#" + filename); 268 } 269 } 270 X.addContent(tagElement); 271 } 272 } 273 } 274 275 /** 276 * Add file section in X. 277 * 278 * @param X XML element to be appended 279 * @param conf Configuration object 280 * @param ns XML element namespace 281 */ 282 static void addFileSection(Element X, Configuration conf, Namespace ns) { 283 String filesStr = conf.get(XOozieClient.FILES); 284 addSection(X, ns, filesStr, "file"); 285 } 286 287 /** 288 * Add archive section in X. 289 * 290 * @param X XML element to be appended 291 * @param conf Configuration object 292 * @param ns XML element namespace 293 */ 294 static void addArchiveSection(Element X, Configuration conf, Namespace ns) { 295 String archivesStr = conf.get(XOozieClient.ARCHIVES); 296 addSection(X, ns, archivesStr, "archive"); 297 } 298}