如何在Shell脚本中判断Spark作业是否成功完成?
没问题,我来给你梳理下怎么实现这个需求——用Shell脚本提交Spark作业并根据执行结果分支处理,这在数据流水线里可太常见了!我分两种常用场景给你写具体实现,毕竟Spark的部署模式不同,处理方式也不一样~
场景1:Spark Client模式(直接等待作业完成)
在Client模式下,spark-submit命令会一直阻塞直到作业执行完毕,并且会返回对应的退出状态码(0代表成功,非0代表失败),这种情况处理起来最简单:
#!/bin/bash # 定义Spark作业的完整提交命令 SPARK_SUBMIT_CMD="spark-submit \ --class com.yourcompany.YourSparkMainClass \ --master yarn \ --deploy-mode client \ --executor-memory 4G \ --num-executors 5 \ hdfs:///path/to/your/spark-job.jar \ --input hdfs:///user/data/input \ --output hdfs:///user/data/output" # 提交作业并输出日志到文件(方便后续排查) echo "启动Spark作业..." $SPARK_SUBMIT_CMD > spark_job_exec.log 2>&1 # 捕获作业的退出状态码 SPARK_EXIT_CODE=$? # 根据状态码分支处理 if [ $SPARK_EXIT_CODE -eq 0 ]; then echo "Spark作业执行成功!开始处理HDFS输出数据..." # 这里写你的后续HDFS操作,示例如下: # 1. 验证输出目录是否存在且有数据 if hdfs dfs -test -d hdfs:///user/data/output && \ [ $(hdfs dfs -count hdfs:///user/data/output | awk '{print $2}') -gt 0 ]; then # 2. 合并小文件(适合分区输出的场景) hdfs dfs -getmerge hdfs:///user/data/output hdfs:///user/data/merged/result.csv # 3. 修改输出目录权限,方便后续访问 hdfs dfs -chmod -R 755 hdfs:///user/data/output echo "后续HDFS操作执行完成!" else echo "警告:输出目录存在但无数据,终止后续操作!" exit 1 fi else echo "Spark作业执行失败,退出码:$SPARK_EXIT_CODE" echo "可查看日志文件 spark_job_exec.log 排查问题" exit $SPARK_EXIT_CODE fi
场景2:YARN Cluster模式(跟踪YARN应用状态)
如果用Cluster模式,默认情况下spark-submit提交后会立即返回(不会等待作业完成),这时候需要通过YARN命令来跟踪作业状态:
#!/bin/bash # 定义Spark作业提交命令 SPARK_SUBMIT_CMD="spark-submit \ --class com.yourcompany.YourSparkMainClass \ --master yarn \ --deploy-mode cluster \ --executor-memory 4G \ --num-executors 5 \ hdfs:///path/to/your/spark-job.jar \ --input hdfs:///user/data/input \ --output hdfs:///user/data/output" # 提交作业并提取YARN应用ID echo "提交Spark作业到YARN..." APP_ID=$($SPARK_SUBMIT_CMD 2>&1 | grep -oP 'application_\d+_\d+') # 检查是否成功获取应用ID if [ -z "$APP_ID" ]; then echo "提交失败:无法获取YARN应用ID,请检查Spark配置或Jar包路径!" exit 1 fi echo "作业提交成功,YARN应用ID:$APP_ID" # 等待作业执行完成 echo "等待作业执行完毕..." yarn application -wait $APP_ID # 获取作业最终状态 APP_STATUS=$(yarn application -status $APP_ID 2>&1 | grep -oP 'Final-State: \K\w+') # 根据状态分支处理 if [ "$APP_STATUS" == "SUCCEEDED" ]; then echo "Spark作业执行成功!开始处理HDFS输出数据..." # 这里写你的后续HDFS操作,和Client模式的示例一致 hdfs dfs -chmod -R 755 hdfs:///user/data/output echo "后续HDFS操作完成!" else echo "Spark作业执行失败,最终状态:$APP_STATUS" echo "可执行命令查看日志:yarn logs -applicationId $APP_ID" exit 1 fi
关键细节与最佳实践
- 日志捕获:把Spark作业的输出/错误日志重定向到文件(比如
> spark_job_exec.log 2>&1),方便后续排查问题。 - 幂等性设计:后续HDFS操作尽量保证幂等,比如执行前先检查输出目录是否存在,避免重复处理导致数据混乱。
- 错误处理:不仅要判断Spark作业的状态,后续HDFS操作也要加入错误检查(比如用
set -e让脚本在命令失败时自动退出,或者手动判断每个命令的退出码)。 - YARN应用ID提取:不同Spark版本的
spark-submit输出格式可能略有不同,如果grep正则匹配不到,可以调整正则表达式或者查看你的Spark版本输出格式。
内容的提问来源于stack exchange,提问作者Manoj Kumar Dhakad
相关产品推荐
相关产品推荐

