Oozie中通过Shell Action捕获Spark作业输出失败问题排查
问题分析与解决方案
我来帮你捋清楚问题出在哪——核心原因是Spark作业的println输出并没有直接流入Shell脚本的标准输出,而Oozie的<capture-output>只捕获Shell脚本本身的stdout内容,并且只会解析key=value格式的行。
当你用spark-submit在Shell脚本里提交Spark作业时,默认情况下,Spark Driver的println内容会被YARN的日志系统接管,不会直接输出到当前Shell的stdout里,自然就被<capture-output>忽略了;而你在Shell里直接echo的内容是直接进入Shell stdout的,所以能被正常捕获。
解决步骤
1. 把Spark的输出重定向到Shell脚本的stdout
修改你的Shell脚本,将Spark提交命令的输出(包括stdout和stderr)重定向到当前Shell的输出,再过滤出你需要的result=xxx行:
# 提交Spark作业,并将所有输出合并到stdout,然后过滤出result开头的行 spark-submit --com.your.package.YourMainClass your-spark-job.jar 2>&1 | grep "result="
这样处理后,Shell脚本的stdout就会只包含result=this is a pen这一行,<capture-output>就能正确解析出result这个变量。
2. 修正Oozie变量读取的键名
你之前用${wf:actionData("action-id")["Variable"]}是错的——因为你的输出是result=xxx,对应的键名是result,不是Variable。后续Action里应该这样读取:
${wf:actionData("your-shell-action-id")["result"]}
备选方案:通过文件中转输出
如果Spark的日志输出太杂,过滤不方便,也可以让Spark直接把结果写入HDFS文件,再让Shell脚本读取并输出成key=value格式:
- Spark代码里写入结果到HDFS:
import org.apache.hadoop.fs.{FileSystem, Path} val fs = FileSystem.get(sc.hadoopConfiguration) val outputFile = new Path("/tmp/spark_result.txt") val writer = fs.create(outputFile) writer.writeBytes("this is a pen") writer.close()
- Shell脚本里读取并格式化输出:
# 读取HDFS文件内容,输出成result=xxx格式 hdfs dfs -cat /tmp/spark_result.txt | awk '{print "result="$0}'
额外注意事项
- 确保Spark的日志级别足够低:如果Spark Driver的日志级别设为
WARN或更高,println的内容可能会被日志系统过滤,需要在spark-submit里加上--conf spark.driver.logLevel=INFO来保证输出能被捕获。 - 避免使用本地路径:如果是YARN集群模式,容器内的本地文件会在作业结束后销毁,所以建议用HDFS路径中转输出。
内容的提问来源于stack exchange,提问作者yAsH
相关产品推荐
相关产品推荐

