Oozie工作流如何基于Scala类内部变量控制执行流程
实现可行性
这个需求完全可以实现,基于Oozie原生的Action输出捕获+决策节点能力即可完成,之前配置的执行报错跳killJobAction的逻辑不会受任何影响。
具体操作步骤
- 改Scala Spark业务代码,输出待判断的变量
Oozie没法直接读取Spark进程内存里的变量,需要你在业务逻辑跑完拿到varWF的取值后,把这个值以键值对形式写到环境变量oozie.action.output.properties指定的HDFS路径下的properties文件里,参考实现:
import java.io.OutputStream
import java.util.Properties
import org.apache.hadoop.fs.{FileSystem, Path}
// 这部分是你原有的业务逻辑,跑完得到varWF的布尔值
val varWF: Boolean = ???
// 把变量写入Oozie识别的输出文件
val hadoopConf = spark.sparkContext.hadoopConfiguration
val fs = FileSystem.get(hadoopConf)
val outputPath = new Path(sys.env("oozie.action.output.properties"))
val outputStream: OutputStream = fs.create(outputPath)
val prop = new Properties()
prop.setProperty("varWF", varWF.toString)
prop.store(outputStream, "workflow transfer params")
outputStream.close()
- 改对应Spark Action的配置,开启输出捕获 在Spark Action的`</spark>`闭合标签之后、`<ok>`标签之前,加一行`<capture-output/>`,告诉Oozie要抓取这个Action输出的属性,同时把这个Action的ok跳转目标改成后续的决策节点,配置示例: ```xml <action name="action name 1" cred="hcat,hs2-creds"> <spark xmlns="uri:oozie:spark-action:0.2"> <!-- 这里原有Spark配置全部保留,不用动 --> <job-tracker>${jobTracker}</job-tracker> <name-node>${nameNode}</name-node> <master>${master}</master> <mode>cluster</mode> <name>class 1 name</name> <class>com.sample.project</class> <jar>${wf_path}/jar_file.jar</jar> <spark-opts>--queue ${queueName} --num-executors 2 --executor-cores 2 --executor-memory 2G --jars ${hiveWarehouseConnectorLib}</spark-opts> </spark> <!-- 新增下面这行,开启输出捕获 --> <capture-output/> <ok to="varWFBranchDecision"/> <error to="killJobAction"/> </action>
- 新增决策节点,配置分支规则
在Spark Action和后续业务节点之间插入decision决策节点,用Oozie内置EL表达式读取抓到的varWF值,匹配对应分支跳转:
踩坑提醒
- Action捕获到的输出值全是字符串类型,写EL判断规则的时候必须和字符串
'true'/'false'比对,不能直接写布尔值判断,会不生效 - 一定要保证Spark任务的执行账号对
oozie.action.output.properties对应的HDFS路径有写权限,否则写输出文件时会直接报错,触发跳转到killJobAction的逻辑 - 如果要传递多个判断变量,直接在properties里多写几个键值对就行,决策节点里可以分别读取做更复杂的分支判断
内容的提问来源于stack exchange,提问作者Sadegh
相关产品推荐
相关产品推荐

